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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 05:02:45