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
相关产品推荐
相关产品推荐

