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

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
原方案问题根源
  1. Chord结果传递机制导致消息超标
    Celery的chord会把group内所有任务的结果收集起来,作为参数传给callback任务。当group_size增大时,start_times列表长度随之增加,序列化后的大小会超过RabbitMQ默认的128MB消息上限。即使group_size=5、chain_size=8,多个chord嵌套在chain中时,构建任务链的过程会生成包含大量任务元数据的消息,同样会触发大小限制。

  2. 任务链构建时的消息膨胀
    在test_chain_of_chords_task中,直接生成多个chord并打包成chain后调用delay(),整个任务链的结构会被序列化发送到RabbitMQ。每个chord包含group的所有任务定义,chain_size较大时,序列化后的任务链结构会异常庞大,直接触发消息大小限制。

  3. 资源配置与并发不匹配
    Celery worker设置了500的并发数,但CPU限制仅为1核,eventlet协程池在CPU资源不足时会出现调度延迟,导致任务堆积,间接加剧消息队列压力,让消息大小问题更早暴露。

最佳实践
  1. 避免Chord直接传递大量结果

    • 不让chord的callback直接接收所有子任务结果,改用结果后端存储数据,让callback任务从Redis中主动获取所需内容(比如只需要第一个任务的开始时间,就只存储和读取该值,而非全部结果)。
    • 子任务直接将关键数据写入数据库或缓存,callback任务仅做触发或简单聚合,不传递大体积数据。
  2. 拆分任务链,避免一次性发送大型任务结构

    • 不在单个任务中构建完整的chain+chord结构,改用递归或中间任务逐步触发后续流程。比如先启动第一个chord,在其callback中触发下一个chord,直到完成所有chain_size次循环,这样每次发送的消息仅包含单个chord的结构,体积可控。
  3. 优化任务结果配置

    • 对不需要返回结果的任务明确设置task_ignore_result=True,仅保留chord中真正需要的任务结果(比如test_request_task如果只需要第一个任务的开始时间,可让其他任务不返回结果)。
    • 缩短result_expires时间,避免结果后端堆积过多无用数据。
  4. 匹配资源与并发数

    • 根据CPU核心数调整Celery worker的并发数,1核CPU下建议将eventlet协程池的并发数控制在100以内,避免调度开销过大。
    • 监控RabbitMQ的内存和消息队列长度,确保broker资源能支撑任务流量。
  5. 使用任务路由与优先级

    • 将不同类型的任务分配到不同的RabbitMQ队列,避免大任务阻塞小任务执行。
    • 为关键任务设置更高优先级,确保流程核心环节不受影响。

内容的提问来源于stack exchange,提问作者Mehedee Ahmed Siddique

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 03:35:04