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

如何在基于Redis的Celery中根据任务ID修改任务结果?

修改Celery任务结果的正确实现方式

一、任务运行中修改结果/状态

如果任务还在执行阶段,有两种可靠方式修改结果:

1. 任务内部主动更新

在任务函数中直接调用self.store_result()覆盖结果,或用self.update_state()更新状态元数据:

from celery import Celery
import time

app = Celery('tasks', backend='redis://localhost:6379/0', broker='redis://localhost:6379/0')

@app.task(bind=True)
def long_running_task(self):
    # 模拟任务执行逻辑
    time.sleep(3)
    # 直接覆盖任务结果
    self.store_result(
        task_id=self.request.id,
        result="Updated task result",
        state="SUCCESS"
    )
    # 这里的return值会被store_result覆盖
    return "Original result"

store_result会直接写入结果后端,优先级高于任务的return值。

2. 外部触发任务内部更新

若要从外部控制正在运行的任务,需要任务内部监听外部指令(比如通过Redis队列):

import redis
from celery import Celery
import time

app = Celery('tasks', backend='redis://localhost:6379/0', broker='redis://localhost:6379/0')
r = redis.Redis(host='localhost', port=6379, db=0)

@app.task(bind=True)
def listenable_task(self):
    while True:
        # 检查是否有更新指令
        update_signal = r.get(f"update_trigger_{self.request.id}")
        if update_signal:
            self.store_result(
                task_id=self.request.id,
                result=update_signal.decode('utf-8'),
                state="SUCCESS"
            )
            r.delete(f"update_trigger_{self.request.id}")
            return
        # 执行常规任务逻辑
        time.sleep(1)

外部触发更新的代码:

r.set("update_trigger_<你的任务ID>", "外部指定的新结果")

二、任务已完成后修改结果

任务执行完毕后,Worker不会再处理该任务的请求,需直接操作结果后端:

Redis后端示例

Celery在Redis中存储任务结果的键格式为celery-task-meta-<任务ID>,直接修改该键即可:

import redis
import json

r = redis.Redis(host='localhost', port=6379, db=0)
task_id = "你的任务ID"

# 构造符合Celery格式的结果数据
new_result = {
    "status": "SUCCESS",
    "result": "已完成任务的更新结果",
    "traceback": None,
    "children": [],
    "date_done": "2024-05-20T15:30:00.123456"
}

# 写入Redis
r.set(f"celery-task-meta-{task_id}", json.dumps(new_result))

SQLAlchemy后端示例

直接修改数据库中celery_taskmeta表的对应记录:

from sqlalchemy import create_engine, update, Column, Integer, String, DateTime, Text, PickleType
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import sessionmaker
import datetime

engine = create_engine('sqlite:///celery_results.db')
Session = sessionmaker(bind=engine)
session = Session()

Base = declarative_base()

# 匹配Celery默认的结果表结构
class CeleryTaskMeta(Base):
    __tablename__ = 'celery_taskmeta'
    id = Column(Integer, primary_key=True)
    task_id = Column(String(255), unique=True)
    status = Column(String(50))
    result = Column(PickleType)
    date_done = Column(DateTime)
    traceback = Column(Text)
    children = Column(PickleType)

# 更新结果
session.execute(
    update(CeleryTaskMeta)
    .where(CeleryTaskMeta.task_id == "你的任务ID")
    .values(
        status="SUCCESS",
        result="已完成任务的更新结果",
        date_done=datetime.datetime.now()
    )
)
session.commit()

常见无效原因排查

  • AsyncResult.update_state无效:通常是任务已完成,或任务未在运行中(该方法是给Worker发消息,仅运行中的任务会处理)。
  • store_result无效:检查是否传入了正确的task_id,以及结果后端的配置、写入权限是否正常。

内容的提问来源于stack exchange,提问作者Juan Camilo R

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 12:47:48