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

如何在Celery中启用任务优先级并获取全部调度任务信息

单队列单并发Celery系统:优先级任务与任务列表查询兼容方案

场景与核心矛盾

我搭建了一套单队列、单并发的Celery任务系统:

  • 1个Celery Worker(并发数1)从RabbitMQ的celery_q队列获取任务
  • Celery-beat定期批量创建低优先级任务(每次800个,优先级1)
  • FastAPI提供两个核心功能:创建最高优先级任务(优先级9)、通过inspect()查询活跃任务进度和剩余待执行任务数

任务代码为循环更新状态的类型:

@c_app.task(bind=True)
def task(self):
    n = 20
    for i in range(0, n):
        self.update_state(state='PROGRESS', meta={'done': i, 'total': n})
        print('working')
        time.sleep(1)
    return n

遇到的核心问题:

  • 初始配置下任务优先级不生效,高优先级任务无法插队执行
  • 添加worker_prefetch_multiplier = 1和task_acks_late = True后优先级生效,但inspect().scheduled()只能看到1个调度任务,无法获取全部待执行任务列表

解决方案

要同时满足高优先级任务插队和完整任务列表查询,需调整Celery配置、任务提交逻辑,并补充RabbitMQ原生查询方式:

1. 调整Celery核心配置

保留优先级生效的关键配置,同时优化队列参数适配大量任务场景:

from celery.schedules import crontab
from kombu import Queue, Exchange

broker_url = 'amqp://guest:guest@localhost//'
result_backend = 'db+postgresql://admin:root@localhost/celery_test_db'

worker_concurrency = 1
timezone = 'Europe/Moscow'
enable_utc = False
result_extended = True

# 优先级生效核心配置
worker_prefetch_multiplier = 1  # 仅预取1个任务,避免Worker提前拉取低优先级任务
task_acks_late = True  # 任务执行完成后再确认,确保未完成时任务留在队列
worker_max_tasks_per_child = 1  # 可选:单任务后重启Worker,防止内存泄漏

beat_schedule = {
    'add-5-tasks-every-month': {
        'task': 'celery_app.tasks.add_5_tasks',
        'options': {'queue': 'celery_q'},
        'schedule': 20.0
    },
}

# RabbitMQ队列优先级与性能配置
broker_transport_options = {'queue_order_strategy': 'priority'}
task_queues = (
    Queue(
        "celery_q", 
        Exchange("celery_q"), 
        routing_key="celery_q", 
        queue_arguments={
            'x-max-priority': 9,  # 队列支持0-9级优先级
            'x-queue-mode': 'lazy'  # 懒加载模式,任务存磁盘,适配大量任务场景
        }
    ),
)

2. 修改任务提交逻辑

移除countdown参数,避免任务被放到Worker本地调度列表(导致inspect()无法获取完整列表):

  • Celery-beat批量任务提交:
@c_app.task
def add_5_tasks():
    for _ in range(800):
        # 移除countdown,直接提交到RabbitMQ队列
        task.apply_async(queue='celery_q', priority=1)
  • FastAPI高优先级任务接口:
@f_app.post("/add-task/")
def add_task():
    # 移除countdown,直接提交到RabbitMQ队列
    task_ = task.apply_async(priority=9, queue='celery_q')
    print('Task added with high priority:', task_.id)
    return {'task_id': task_.id,
            'message': 'Task added with high priority'}

3. 调整任务查询逻辑

由于inspect().scheduled()仅能获取Worker本地预取的任务,需通过RabbitMQ原生API查询队列中的任务总数:

import pika

def get_queue_task_count(queue_name='celery_q'):
    connection = pika.BlockingConnection(pika.URLParameters('amqp://guest:guest@localhost//'))
    channel = connection.channel()
    queue_declare = channel.queue_declare(queue=queue_name, passive=True)
    task_count = queue_declare.method.message_count
    connection.close()
    return task_count

# FastAPI查询接口示例
@f_app.get("/task-stats/")
def get_task_stats():
    i = c_app.control.inspect()
    active = i.active() or {}
    active_count = sum(len(tasks) for tasks in active.values())
    # 获取RabbitMQ队列中未被预取的任务数
    queue_task_count = get_queue_task_count()
    return {
        'active_task_count': active_count,
        'queued_task_count': queue_task_count,
        'total_pending_tasks': active_count + queue_task_count
    }

原理说明

  • worker_prefetch_multiplier = 1 + task_acks_late = True:确保Worker仅持有当前正在执行的任务,完成后才从RabbitMQ拉取下一个,此时RabbitMQ会按优先级返回任务,实现高优先级插队
  • 移除countdown:让所有待执行任务留在RabbitMQ队列中,而非Worker本地调度列表,避免inspect()无法获取完整任务列表
  • x-queue-mode: lazy:优化大量任务场景下的RabbitMQ内存占用,避免内存溢出

内容的提问来源于stack exchange,提问作者Mika

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 08:57:12