【问题标题】:Celery's task_reject_on_worker_lost doesn't work with Redis as message brokerCelery 的 task_reject_on_worker_lost 不能使用 Redis 作为消息代理
【发布时间】:2022-12-14 15:11:26
【问题描述】:
我目前使用的是 Celery 5.2.6 版和 Redis 6.2.6 版。当我打开 task_reject_on_worker_lost 标志时,我希望 Celery 重新交付由突然死亡的工人执行的任务。但是,在 Redis 上尝试将此作为消息代理,我的任务实际上并没有在工作人员宕机后立即重新交付。另一方面,当我尝试使用 RabbitMQ 进行完全相同的配置时,它会按预期工作。
关于如何使用 Redis 作为消息代理实现相同行为的任何指示?
【问题讨论】:
标签:
celery
django-celery
celery-task
celeryd
djcelery
【解决方案1】:
我最近刚接触芹菜,面临着和你一样的问题。
这意味着使用 ack 配置:
task_acks_late = True # ack after task end
task_acks_on_failure_or_timeout = True # ack if task exception
task_reject_on_worker_lost = True # no ack if worker killed
如果代理配置使用 redis:
broker_url = f'redis://127.0.0.1:6379/1'
如果 worker 在运行任务期间被杀死并再次重新启动,任务将不会重新排队。
但是如果使用rabbitmq:
broker_url = 'amqp://guest:guest@localhost:5672/'
任务重新排队等待运行。
最后,我从 celery github issues 得到了这个comment。
broker_transport_options 中 visibility_timeout 的附加配置值我需要 redis。
我在我的配置中添加了额外的配置并且它正在工作。
仅供参考,这是我的配置文件:
broker_url = f'redis://127.0.0.1:6379/1'
# task message ack
# https://docs.celeryq.dev/en/stable/userguide/configuration.html#std-setting-task_acks_late
# https://docs.celeryq.dev/en/stable/userguide/configuration.html#task-acks-on-failure-or-timeout
# https://docs.celeryq.dev/en/stable/userguide/configuration.html#task-reject-on-worker-lost
task_acks_late = True # ack after task end
task_acks_on_failure_or_timeout = True # ack if task exception
task_reject_on_worker_lost = True # no ack if worker killed
# only for redis broker
# https://github.com/celery/celery/issues/4984
broker_transport_options = {'visibility_timeout': 10}
import celery
import celery_config
app = celery.Celery("celery")
app.config_from_object(celery_config)