如何在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
相关产品推荐
相关产品推荐

