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

如何将Celery用作优先级队列?Django后台任务需求问询

解决方案

一、用Celery+Redis实现优先级队列(无需自行实现排序)

Celery结合Redis Broker本身就支持成熟的优先级队列机制,比手动实现排序更能规避竞态条件,直接用官方方案即可:

  1. 配置多优先级队列
    在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,
}
  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')
    )
  1. 优化锁机制
    保留原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更直接):

  1. 检测高优先级任务
    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
  1. 任务中断并重入
    在任务执行的检查点,若检测到高优先级任务,终止当前操作(需提前确保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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 21:43:11