持续部署场景下处理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
相关产品推荐
相关产品推荐

