You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Celery 4.4.6如何立即终止运行任务并获取任务运行状态

问题根因

两个方案失效的原因非常明确:

  • app.control.revoke(terminate=True, signal='SIGKILL') 对内嵌asyncio事件循环的任务无效:Celery 4.x的terminate逻辑是给worker拉起的任务执行进程发系统信号,但你的无限循环逻辑跑在asyncio.run()启动的独立事件循环协程里,信号不会透传到协程层,协程一直占着执行流,进程不会直接退出。如果开了任务重试配置,进程被杀后任务还会被worker重新拉起,完全达不到停止效果。
  • Django post_save 回调不触发是常识错误:Django信号是单进程内的内存回调机制,你在Celery worker进程里注册的信号回调,只能收到同一个worker进程内触发的模型save事件;前端点停止触发的save()是跑在Web服务进程里的,两个进程内存完全隔离,worker根本收不到这个事件,回调永远不会执行。
可直接落地的实现方案

1. 任务逻辑改造:轮询停止标记+协程安全退出

废弃跨进程信号、强杀进程的思路,用数据库作为公共状态存储,任务每次循环主动检查停止标记,正常退出前可完成资源清理,无脏数据、无进程泄漏。

import asyncio
from celery import shared_task
from myapp.models import Task, TaskSignal

async def _grind(task_id: str):
    count = 1
    pending_req_tasks = set()
    aiohttp_session = None # 初始化你自己的aiohttp会话
    try:
        while True:
            # 每次循环开头先查停止标记,停止延迟由循环间隔控制,可做到100ms内响应
            task_record = Task.objects.filter(id=task_id).first()
            if task_record and task_record.signal == TaskSignal.STOP:
                # 退出前清理:取消所有未完成的aiohttp请求、关闭会话、提交未入库数据、释放锁
                for t in pending_req_tasks:
                    t.cancel()
                await asyncio.gather(*pending_req_tasks, return_exceptions=True)
                if aiohttp_session and not aiohttp_session.closed:
                    await aiohttp_session.close()
                task_record.delete()
                return

            count += 1
            print(count)
            # 这里放你的业务逻辑:aiohttp请求、数据库写入,注意给单个请求加超时
            # 示例:管理并发请求,避免协程泄漏
            req_task = asyncio.create_task(aiohttp_session.get("your_api_url", timeout=10))
            pending_req_tasks.add(req_task)
            req_task.add_done_callback(pending_req_tasks.discard)

            # 控制最大并发数,按你的业务需求调整
            if len(pending_req_tasks) >= 10:
                done, pending_req_tasks = await asyncio.wait(
                    pending_req_tasks, 
                    return_when=asyncio.FIRST_COMPLETED
                )
                # 处理已完成请求的结果,写入数据库
    except asyncio.CancelledError:
        # 兜底处理协程被意外取消的场景,清理残留数据
        Task.objects.filter(id=task_id).delete()

@shared_task(bind=True, track_started=True)
def grind(self):
    task_id = self.request.id
    asyncio.run(_grind(task_id))

注意:必须给装饰器加track_started=True参数,否则Celery不会记录任务运行中状态,后续状态查询会不准。

2. 启停接口逻辑修正

不需要强杀进程,修改数据库停止标记即可,任务轮询到标记会自动退出:

from myapp.tasks import grind

def toggle_grind(user):
    started_task = Task.objects.filter(type=TaskType.GRIND, user=user).first()
    if started_task:
        started_task.signal = TaskSignal.STOP
        started_task.save()
        # 可选:补一个不带terminate的revoke,把还没轮到执行的任务从队列里移除
        app.control.revoke(started_task.id)
        return False
    else:
        async_result = grind.apply_async()
        Task.objects.create(
            id=async_result.id,
            type=TaskType.GRIND,
            user=user
        )
        return True

3. 任务状态查询接口

结合Celery原生状态和业务表状态判断,自动清理异常退出产生的脏数据:

from celery.result import AsyncResult

def get_grind_status(user):
    task_record = Task.objects.filter(type=TaskType.GRIND, user=user).first()
    if not task_record:
        return {"is_running": False}
    res = AsyncResult(task_record.id)
    # Celery状态说明:PENDING=未调度、STARTED=运行中、SUCCESS/FAILURE/REVOKED=已结束
    if res.state == "STARTED":
        return {"is_running": True, "task_id": task_record.id}
    # 任务已结束但数据库有残留记录,说明是异常退出,清理脏数据
    if res.state in ("SUCCESS", "FAILURE", "REVOKED"):
        task_record.delete()
    return {"is_running": False}

4. 必加的Celery配置

否则状态查询会出现偏差:

# 开启任务启动状态记录
CELERY_TRACK_STARTED = True
# 配置结果存储,推荐用django-celery-results存在业务数据库,和Task表共用数据源
CELERY_RESULT_BACKEND = 'django-db'
# 任务结果过期时间,按需配置,示例为1天
CELERY_RESULT_EXPIRES = 86400
# 关闭任务默认重试,避免异常退出后任务自动拉起
CELERY_ACKS_LATE = False

内容的提问来源于stack exchange,提问作者RomanGodMode

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 03:12:19