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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 20:50:32