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

持续部署场景下处理Celery未注册任务的方案

解决Celery同队列worker任务不兼容的问题

核心方案:拦截未注册任务并重新入队

直接用Celery的on_task_received信号做拦截,让没注册对应任务的worker把任务扔回队列,交给有能力处理的worker。

具体实现

给所有任务不完整的worker(比如机器A)加一段信号处理代码就行:

from celery import Celery
from celery.signals import on_task_received

app = Celery('worker', broker='pyamqp://guest@localhost//', backend='redis://localhost:6389/0')

@app.task
def add(x, y):
    return x + y

@on_task_received.connect(sender=app)
def skip_unknown_tasks(sender, task, **kwargs):
    # 检查当前worker有没有这个任务的注册
    if task.name not in sender.tasks:
        # 拒绝任务并放回队列,让其他worker接手
        task.reject(requeue=True)

机器B的代码不用改,保持原样即可。

为什么这么做能行

  • on_task_received是worker刚拿到任务还没执行时触发的信号,刚好能在报错前拦截。
  • 调用task.reject(requeue=True)不会把任务标记为失败,而是送回原队列,队列会重新分配给其他worker,直到找到能处理的那个。
  • 完全不用拆分队列,完美适配你的场景。

注意点

  • 所有任务不全的worker都得加这段代码,漏加的话还是会抛异常。
  • 任务重新入队会有一点延迟,但不会丢任务,属于可接受的范围。

内容的提问来源于stack exchange,提问作者naktinis

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 21:37:02