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

如何在Celery的celeryd_init信号中读取-Q指定队列并修改Worker并发数

解决Celery Worker按指定队列配置并发数的问题

核心问题原因

celeryd_init信号触发时,Worker尚未处理命令行传入的-Q参数,因此celery_app.amqp.queues只会包含默认队列;而celeryd_after_setup等后续信号触发时,队列信息已加载,但此时Worker的进程池已初始化,直接修改celery_app.conf无法生效。


方案一:在celeryd_init中解析命令行参数获取队列

直接从启动命令的参数中提取队列名称,无需等待App的队列配置加载:

@celeryd_init.connect
def configure_workers(sender=None, **kwargs):
    argv = kwargs.get('argv', [])
    queue_names = []
    
    # 解析命令行中的-Q参数,支持逗号分隔多队列
    try:
        q_flag_index = argv.index('-Q')
        queue_names = argv[q_flag_index + 1].split(',')
    except (ValueError, IndexError):
        # 未指定-Q时使用默认队列
        queue_names = ['celery']

    # 根据队列设置对应并发数
    if 'celery' in queue_names:
        celery_app.config_from_object(config, namespace='CELERY')
        celery_app.conf.update(worker_concurrency=4)
    elif 'queue2' in queue_names:
        celery_app.config_from_object(config, namespace='queue2')
        celery_app.conf.update(worker_concurrency=2)

方案二:在celeryd_after_setup中修改Worker实例并发数并重启进程池

利用celeryd_after_setup信号获取已初始化的Worker实例,直接修改其并发数并重启进程池生效:

@celeryd_after_setup.connect
def setup_worker_concurrency(sender=None, **kwargs):
    # sender即为当前Worker实例
    queue_names = list(sender.app.amqp.queues.keys())
    
    target_concurrency = 4
    if 'queue2' in queue_names:
        target_concurrency = 2
    
    # 修改并发数并重启进程池应用配置
    sender.concurrency = target_concurrency
    sender.pool_restart()

方案说明

  • 方案一:在Worker初始化早期完成配置,无需重启进程池,适合需要加载完整配置文件的场景。
  • 方案二:依赖Worker实例的已加载队列信息,通过重启进程池应用新并发数,逻辑更直观,但会短暂重启Worker进程。

内容的提问来源于stack exchange,提问作者Learning from masters

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 19:40:31