Celery任务调度后Worker与请求方双向通信实现方法问询
问题:Celery中复用现有连接实现Worker与任务请求方的双向通信
我已经通过update_state和on_message回调实现了Celery中Worker到任务请求方的单向进度更新通信,但现在需要在调用long_task.apply_async()之后,给Worker发送后续才获取到的任务额外信息。想知道能不能复用现有的连接实现这种双向通信,不用新建连接。
代码示例:
@app.task(bind=True) def long_task(self: celery.Task): time.sleep(1) self.update_state(state="PROGRESS", meta={"progress": 50}) time.sleep(1) self.update_state(state="PROGRESS", meta={"progress": 90}) time.sleep(1) self.update_state(state="DONE", meta={"progress": 100}) return "done" result: celery.result.AsyncResult = long_task.apply_async() result.get(on_message=on_message) # on_message回调会接收进度更新
回答
可以复用现有连接实现双向通信,核心是基于Celery已有的Broker连接,通过内置机制传递额外信息,以下是两种可行方案:
方案一:用
send_task发送指令任务
直接通过现有Celery实例的send_task方法发送一个"指令任务",该方法会复用已建立的Broker连接,无需新建连接。
请求方发送额外信息的代码:# 后续获取到额外信息后执行 app.send_task('long_task.receive_extra_data', args=[result.id, extra_info])Worker端改造任务逻辑,接收并存储额外信息:
# 单进程环境用全局字典存活跃任务,多进程/多Worker需改用Redis等共享存储 active_tasks = {} @app.task(bind=True) def long_task(self: celery.Task): active_tasks[self.request.id] = self try: time.sleep(1) self.update_state(state="PROGRESS", meta={"progress": 50}) # 后续任务逻辑可以直接使用self.extra_data time.sleep(1) self.update_state(state="PROGRESS", meta={"progress": 90}) time.sleep(1) self.update_state(state="DONE", meta={"progress": 100}) return "done" finally: del active_tasks[self.request.id] @app.task def receive_extra_data(task_id, extra_info): task = active_tasks.get(task_id) if task: task.extra_data = extra_info方案二:反向利用
update_state传递信息
请求方调用AsyncResult.update()方法向Broker发送自定义状态消息,Worker通过监听任务状态变化获取额外信息,全程复用现有连接:
请求方代码:# 后续获取到额外信息后,发送自定义状态消息 result.update(state="EXTRA_DATA", meta={"data": extra_info})Worker端监听状态变化:
from celery.signals import task_state_changed active_tasks = {} @task_state_changed.connect(sender='long_task') def handle_extra_data(sender, task_id, state, meta, **kwargs): if state == "EXTRA_DATA" and task_id in active_tasks: active_tasks[task_id].extra_data = meta['data'] @app.task(bind=True) def long_task(self: celery.Task): active_tasks[self.request.id] = self try: # 原有任务逻辑... finally: del active_tasks[self.request.id]关键注意事项
- 多Worker/多进程环境下,不能用全局字典存储活跃任务,必须改用Redis、Memcached等共享存储来跨进程访问任务实例。
- 自定义状态名称(如方案二中的
EXTRA_DATA)要避免和Celery默认状态(PROGRESS、DONE等)冲突。 - 两种方案均基于Celery原有Broker连接,无需新建额外连接,完全满足你的需求。
内容的提问来源于stack exchange,提问作者Piotr Dabkowski
相关产品推荐
相关产品推荐

