Celery任务流触发RabbitMQ消息大小超限问题排查与优化咨询
我设计了包含test_request_task、test_chord_callback_task、test_chain_of_chord_callback_task和test_chain_of_chords_task的Celery任务流,需求为每秒执行n次test_request_task(1≤n≤100),持续m秒(5≤m≤30)。采用Celery的group、chain、chord构建该流程,当m和n取值较小时无异常,但参数增大后任务执行变慢,最终触发错误:
PreconditionFailed(406, 'PRECONDITION_FAILED - message size is larger than configured max size 134217728', (60, 40), 'Basic.publish')
我已通过拆分流程并委托给中间任务解决了问题,但想明确原方案的问题所在,同时寻求相关最佳实践。由于实际payload可能动态增大,调大RabbitMQ默认消息大小并非理想方案。
我使用RabbitMQ作为broker,Redis作为结果后端(仅保留chord内某一任务的结果),相关配置与任务代码如下:
CELERY_BROKER: str = ( f"pyamqp://xxx:Zzzz1234@rabbitmq:5672//" ) CELERY_BACKEND: str = f"redis://redis:6379" celery_app = Celery( __name__, broker=CELERY_BROKER, backend=CELERY_BACKEND ) celery_app.conf.broker_pool_limit = 500 celery_app.conf.worker_prefetch_multiplier = 5 celery_app.conf.worker_send_task_events = False celery_app.conf.task_send_sent_event = False celery_app.conf.task_ignore_result = True celery_app.conf.task_store_errors_even_if_ignored = True celery_app.conf.result_expires = 30 @celery_app.task(ignore_result=False) def test_request_task(serial: int): """ The unit task. :param serial: int :return: starting time of this task """ logger.info(f"Task {serial} started") start_time = time.time() TestModel.objects.create(req_number=serial, req_exec=True) try: resp = requests.get(f"http://mock_app:10034/test_request?serial={serial}") TestModel(req_number=serial).update(rec_ack=True) except Exception as e: logger.error(f"Task {serial} - ERROR: {e.__str__()}") return start_time @celery_app.task def test_chord_callback_task(start_times: list[float]): """ Task to ensure I don't execute more than n test_request_task in 1 sec. If all test_request_task complete under 1 sec then sleep for the remainder of that sec, Otherwise do nothing. :param start_times: list[float] """ now = time.time() logger.info(f"Group tasks finished") time_since_first_task_in_grp = now - start_times[0] if len(start_times) > 0 else 0 if time_since_first_task_in_grp < 1.0: sleep_time_sec = 1.0 - time_since_first_task_in_grp logger.info(f"SLEEPING {sleep_time_sec}") time.sleep(sleep_time_sec) @celery_app.task def test_chain_of_chord_callback_task(): """ Just another task that runs after all tasks have finished executing :return: """ logger.info("Chain of chord finished") @celery_app.task def test_chain_of_chords_task(group_size: int, chain_size: int): """ Make chord with a group of group_size test_request_task to run in parallel and ensure 1 sec completion, with callback test_chord_callback_task. Make chain_size numbers of above chords and chain them with additional test_chain_of_chord_callback_task. :param group_size: :param chain_size: """ total_size = group_size * chain_size chords = ( chord( group( test_request_task.si(serial) for serial in range(i, i + group_size) ), test_chord_callback_task.s() ) for i in range(0, total_size, group_size) ) chain_of_chords_with_callback = chain( *chords, test_chain_of_chord_callback_task.si() ) logger.info("Firing chain of chords") chain_of_chords_with_callback.delay()
即使设置group_size=5、chain_size=8,该流程仍会崩溃。Docker-Compose中的资源分配如下:
celery_worker: restart: always build: . volumes: - .:/code command: - "celery" - "-A" - "app.worker.celery_app" - "worker" - "--pool=eventlet" - "--concurrency=500" - "--loglevel=INFO" depends_on: - rabbitmq - redis deploy: resources: limits: cpus: '1.0' memory: 2G reservations: cpus: '0.50' memory: 500M rabbitmq: image: rabbitmq:3.11.3-management-alpine ports: - "5672:5672" - "15672:15672" environment: - RABBITMQ_HOST=rabbitmq - RABBITMQ_PORT=5672 - RABBITMQ_DEFAULT_USER=xxx - RABBITMQ_DEFAULT_PASS=Zzzz1234 deploy: resources: limits: cpus: '0.5' memory: 1G reservations: cpus: '0.50' memory: 500M redis: image: redis:7.0.5-alpine ports: - "6379:6379" environment: - REDIS_HOST=redis - REDIS_PORT=6379 deploy: resources: limits: cpus: '0.5' memory: 100M reservations: cpus: '0.50' memory: 50M
Chord结果传递机制导致消息超标
Celery的chord会把group内所有任务的结果收集起来,作为参数传给callback任务。当group_size增大时,start_times列表长度随之增加,序列化后的大小会超过RabbitMQ默认的128MB消息上限。即使group_size=5、chain_size=8,多个chord嵌套在chain中时,构建任务链的过程会生成包含大量任务元数据的消息,同样会触发大小限制。任务链构建时的消息膨胀
在test_chain_of_chords_task中,直接生成多个chord并打包成chain后调用delay(),整个任务链的结构会被序列化发送到RabbitMQ。每个chord包含group的所有任务定义,chain_size较大时,序列化后的任务链结构会异常庞大,直接触发消息大小限制。资源配置与并发不匹配
Celery worker设置了500的并发数,但CPU限制仅为1核,eventlet协程池在CPU资源不足时会出现调度延迟,导致任务堆积,间接加剧消息队列压力,让消息大小问题更早暴露。
避免Chord直接传递大量结果
- 不让chord的callback直接接收所有子任务结果,改用结果后端存储数据,让callback任务从Redis中主动获取所需内容(比如只需要第一个任务的开始时间,就只存储和读取该值,而非全部结果)。
- 子任务直接将关键数据写入数据库或缓存,callback任务仅做触发或简单聚合,不传递大体积数据。
拆分任务链,避免一次性发送大型任务结构
- 不在单个任务中构建完整的chain+chord结构,改用递归或中间任务逐步触发后续流程。比如先启动第一个chord,在其callback中触发下一个chord,直到完成所有chain_size次循环,这样每次发送的消息仅包含单个chord的结构,体积可控。
优化任务结果配置
- 对不需要返回结果的任务明确设置
task_ignore_result=True,仅保留chord中真正需要的任务结果(比如test_request_task如果只需要第一个任务的开始时间,可让其他任务不返回结果)。 - 缩短
result_expires时间,避免结果后端堆积过多无用数据。
- 对不需要返回结果的任务明确设置
匹配资源与并发数
- 根据CPU核心数调整Celery worker的并发数,1核CPU下建议将eventlet协程池的并发数控制在100以内,避免调度开销过大。
- 监控RabbitMQ的内存和消息队列长度,确保broker资源能支撑任务流量。
使用任务路由与优先级
- 将不同类型的任务分配到不同的RabbitMQ队列,避免大任务阻塞小任务执行。
- 为关键任务设置更高优先级,确保流程核心环节不受影响。
内容的提问来源于stack exchange,提问作者Mehedee Ahmed Siddique

