如何将Celery用作优先级队列?Django后台任务需求问询
解决方案
一、用Celery+Redis实现优先级队列(无需自行实现排序)
Celery结合Redis Broker本身就支持成熟的优先级队列机制,比手动实现排序更能规避竞态条件,直接用官方方案即可:
- 配置多优先级队列
在Django的settings.py中定义不同优先级的队列,给高优先级队列设置更高权重,让Worker优先消费:
CELERY_TASK_QUEUES = { 'high_priority': { 'exchange': 'high_priority', 'routing_key': 'high_priority', }, 'normal_priority': { 'exchange': 'normal_priority', 'routing_key': 'normal_priority', }, 'low_priority': { 'exchange': 'low_priority', 'routing_key': 'low_priority', }, } # 配置Worker消费权重:数值越高,队列优先级越高 CELERY_WORKER_QUEUES = { 'high_priority': 3, 'normal_priority': 2, 'low_priority': 1, }
- 发送任务时指定队列
在视图中根据业务逻辑选择对应优先级的队列发送任务:
def some_view(request): priority = request.GET.get('priority', 'normal') queue_map = { 'high': 'high_priority', 'normal': 'normal_priority', 'low': 'low_priority' } my_task.apply_async( args=[request.GET.get('data1'), request.GET.get('data2')], kwargs={'priority': priority}, queue=queue_map.get(priority, 'normal_priority') )
- 优化锁机制
保留原Redis锁的同时,设置合理超时时间,避免任务意外挂掉后锁长期占用:
@shared_task(bind=True, max_retries=3) def my_task(self, data1, data2, priority='normal'): # 锁超时设为300秒,可根据API实际调用耗时调整 with cache.lock("my_task", timeout=300): use_api(data1) # 插入高优先级任务检查逻辑,下文详细说明 check_and_requeue(self, priority) use_api(data2)
二、任务中途检查高优先级任务并重新入队
要实现运行中任务主动检查队列并重入,直接操作Redis队列效率最高(Celery的Inspect API适合集群场景,但单实例Redis更直接):
- 检测高优先级任务
Redis中Celery队列以列表形式存储(键名格式为celery:<queue_name>),读取队列头部任务并解析优先级:
import json from redis import Redis def has_high_priority_tasks(): redis_client = Redis() high_queue_key = 'celery:high_priority' # 只检查队列前10个任务,避免遍历过长队列影响性能 tasks = redis_client.lrange(high_queue_key, 0, 10) for task in tasks: try: task_data = json.loads(task) # 从任务参数中读取优先级 if task_data.get('kwargs', {}).get('priority') == 'high': return True except json.JSONDecodeError: continue return False
- 任务中断并重入
在任务执行的检查点,若检测到高优先级任务,终止当前操作(需提前确保use_api支持中断,比如给HTTP请求设超时、关闭连接),然后将当前任务重新入队:
def check_and_requeue(task_instance, current_priority): # 仅非高优先级任务需要检查 if current_priority != 'high' and has_high_priority_tasks(): # 这里根据实际场景清理资源,比如关闭API连接 # 将当前任务延迟5秒重入队,给高优先级任务让路 task_instance.retry( args=task_instance.args, kwargs=task_instance.kwargs, queue=f"{current_priority}_priority", countdown=5 )
修改后的完整代码
from celery import shared_task from django.core.cache import cache import json from redis import Redis @shared_task(bind=True, max_retries=3) def my_task(self, data1, data2, priority='normal'): with cache.lock("my_task", timeout=300): use_api(data1) # 检查高优先级任务并决定是否重入队 if current_priority != 'high' and has_high_priority_tasks(): # 清理当前API相关资源(示例,需根据实际情况调整) # close_api_connection() self.retry( args=[data1, data2], kwargs={'priority': priority}, queue=f"{priority}_priority", countdown=5 ) use_api(data2) def has_high_priority_tasks(): redis_client = Redis() high_queue_key = 'celery:high_priority' tasks = redis_client.lrange(high_queue_key, 0, 10) for task in tasks: try: task_data = json.loads(task) if task_data.get('kwargs', {}).get('priority') == 'high': return True except json.JSONDecodeError: continue return False def some_view(request): priority = request.GET.get('priority', 'normal') queue_map = { 'high': 'high_priority', 'normal': 'normal_priority', 'low': 'low_priority' } my_task.apply_async( args=[request.GET.get('data1'), request.GET.get('data2')], kwargs={'priority': priority}, queue=queue_map.get(priority, 'normal_priority') )
内容的提问来源于stack exchange,提问作者allo
相关产品推荐
相关产品推荐

