Celery current_task.update_state不生效问题排查与解决
DRF 结合 Celery+RabbitMQ 实现异步任务进度更新的问题解决
问题描述
在Django REST Framework服务中,通过RabbitMQ作为消息中间件、Celery Worker实现长耗时异步任务,需求是间隔更新任务进度状态。目前任务提交和最终结果获取均正常,但在任务执行过程中及完成前后查询状态时,始终返回task.result = None、task.state = PENDING、task.info = None。
项目结构
DeployML/ manage.py DeployML/ ... settings.py tasks.py dcelery.py views.py urls.py
相关代码片段
dcelery.py 代码
import os from celery import Celery import sys sys.path.append('.....DeployML') sys.path.append('.....DeployML/DeployML') # Set the default Django settings module for the 'celery' program. os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'DeployML.settings') app = Celery('DeployML', backend='rpc://', broker='pyamqp://guest@localhost//', worker_state_db = '/tmp/celery_state') app.config_from_object('django.conf:settings') app.autodiscover_tasks() @app.task(bind=True) def debug_task(self): print(f'Request: {self.request!r}')
tasks.py 代码
import os os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'DeployML.settings') from celery import shared_task from celery import current_task from dcelery import app import time @app.task(bind=True) def addShrd(self, x, y): for i in range(4): time.sleep(5) progress_percent = i #NOTE update_state not working yet current_task.update_state(state='PROGRESS', meta={'current': i, 'total': 4}) return x + y
views.py 代码
#... #... @api_view(['GET', 'POST']) def celer_view(request): print("authentication: ", request.user, request.auth) if request.method == 'GET': result = addShrd.delay(17,10) return Response({'task_id': str(result.id)}) elif request.method == 'POST': return Response({'some': 'POST response data'}) return Response({'some': 'data'}) from celery.result import AsyncResult @api_view(['GET', 'POST']) def poll_state(request): data = 'Fail' print("requ: "+str(request.POST.get('task_id'))) if 'task_id' in request.POST.keys() and request.POST['task_id']: task_id = request.POST['task_id'] task = AsyncResult(task_id) print("meta data: "+ str(task.info)) data = str(task.result) + " " +str(task.state) + " | " + str(task.info) else: data = 'No task_id in the request' return Response({'some': 'poll_state data: '+data})
问题排查与解决
- 排查尝试:曾将
poll_state方法中AsyncResult的导入从dcelery.AsyncResult改为celery.result.AsyncResult,问题未解决。 - 最终解决:发现
update_state更新状态与AsyncResult获取状态之间存在几秒延迟,原任务总时长仅20秒,RabbitMQ还未完成中间进度状态的持久化,任务就已执行完毕。将任务时长延长至100秒后,可正常查询到中间进度结果。
内容的提问来源于stack exchange,提问作者Kshitiz Sharma
相关产品推荐
相关产品推荐

