Celery多消费者重复消费同一RabbitMQ任务问题排查
问题
在一台服务器部署RabbitMQ后,6台机器通过Celery连接该RabbitMQ,发现同一任务ID(ef00cc1f-1be5-44ba-8911-90c0746196ba)的tasks.task_name任务被多台机器重复接收,求排查原因及缺失配置。
现有配置
Celery配置
import os from celery import Celery import settings import sys sys.path.insert(0, f"{settings.BASE_DIR.parent}/") broker_url = "amqp://username:password@machine_ip:rabbitmq_port/" app = Celery(broker=broker_url) app.conf.update( task_acks_late=True, broker_transport_options={'visibility_timeout': 3600}, # Adjust timeout as needed broker_connection_retry_on_startup = True, imports = ['tasks'] ) app.autodiscover_tasks() # celery -A celery_app worker --pool=solo -Q queue_name -l info --logfile logs/celery.logs
RabbitMQ Docker Compose配置
rabbitmq: image: "rabbitmq:3-management" ports: - "5677:5672" - "15677:15672" environment: RABBITMQ_DEFAULT_USER: "username" RABBITMQ_DEFAULT_PASS: "password" RABBITMQ_DEFAULT_VHOST: "/" volumes: - ./rabbitmq_data:/var/lib/rabbitmq
Celery任务代码
@app.task(queue="queue_name") def task_name(batch_s3_path, task_id, api_data, index): pass
问题原因与解决方案
核心原因
当前配置未启用任务ID唯一去重机制,Celery默认不会基于任务ID阻止重复分发。当同一ID的任务被多次发布、或因连接/Worker故障导致RabbitMQ重新投递时,多个Worker会同时接收该任务。
缺失的关键配置
需添加以下配置来实现任务ID级别的唯一约束:
- 全局启用任务唯一校验
在Celery配置中添加task_unique_on: True,该配置会让Broker层面基于任务ID拦截重复任务:
app.conf.update( # 保留原有配置 task_unique_on = True, )
- 发布任务时显式指定唯一性
如果是通过apply_async发布任务,需同时指定task_id和unique=True参数,确保同一ID的任务不会被重复入队:
task_name.apply_async( args=(batch_s3_path, task_id, api_data, index), task_id="ef00cc1f-1be5-44ba-8911-90c0746196ba", queue="queue_name", unique=True )
其他排查点
- 检查业务代码是否存在重复发布任务的逻辑,比如循环调用任务发布接口导致同一ID任务多次入队。
- 验证
visibility_timeout是否合理:如果任务实际执行时间超过3600秒,RabbitMQ会将任务重新放回队列,被其他Worker接收,需根据任务耗时调整该值。 - 查看Worker日志:如果Worker执行任务时崩溃且未发送ACK,RabbitMQ会自动重新投递任务,这种情况属于故障恢复,需排查Worker崩溃原因。
内容的提问来源于stack exchange,提问作者mohamed naser
相关产品推荐
相关产品推荐

