如何在基于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
相关产品推荐
相关产品推荐

