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

Celery中start_distribution调用send_request组任务未入指定队列问题

问题描述

在start_distribution任务中通过apply_async调用send_request组任务未生效:start_distribution已成功执行,但提交的组任务并未进入external_service_req队列。


相关代码

tasks.py

from celery import shared_task, group
from celery.utils.log import get_task_logger

from .models import Distribution, Client

logger = get_task_logger(__name__)


@shared_task(bind=True, max_retries=200, default_retry_delay=30)
def send_request(self, distribution_id, client_id):
    pass


@shared_task(bind=True)
def start_distribution(self, distribution_id):

    distribution = Distribution.objects.get(pk=distribution_id)

    logger.info(f'distribution obj:')
    logger.info(f'  pk: {distribution.pk}')
    logger.info(f'  start date: {distribution.start_date}')
    logger.info(f'  end date: {distribution.end_date}')
    logger.info(f'  tag: {distribution.tag.name}')
    logger.info(f'  mobile_code : {distribution.mobile_code.mobile_code}')
    logger.info(f'  body message : {distribution.body_message}')

    clients_for_send = Client.objects.filter(
        mobile_code=distribution.mobile_code,
        tag=distribution.tag
    ).all()

    logger.info('Pk sending client:')

    [logger.info(client.pk) for client in clients_for_send]

    g = group(send_request.s(distribution_id, client.pk) for client in clients_for_send)

    g.apply_async(
        queue='external_req',
        expires=distribution.end_date
    )

celery.py

from celery import Celery

# Set the default Django settings module for the 'celery' program.
os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'Notification_service.settings')

app = Celery('Notification_service',
             broker='redis://localhost/0',
             backend='redis://localhost/1'
             )

# Using a string here means the worker doesn't have to serialize
# the configuration object to child processes.
# - namespace='CELERY' means all celery-related configuration keys
#   should have a `CELERY_` prefix.
app.config_from_object('django.conf:settings', namespace='CELERY')

app.conf.task_queues = (
    Queue('distribution', routing_key='distribution'),
    Queue('external_service_req', routing_key='external_req')
)

# Load task modules from all registered Django apps.
app.autodiscover_tasks()

Worker启动命令

  • 针对distribution队列:
celery -A NotificationService.celery worker -P prefork -l INFO -Q distribution -n distribution
  • 针对external_req队列:
celery -A NotificationService.celery worker -P gevent -l INFO -Q external_service_req -n external_service_req

日志信息

  • distribution worker日志显示start_distribution任务成功完成
  • external_service_req worker日志仅显示就绪状态,未接收到任务

排查与解决方案

1. 修正队列名称匹配问题

你在g.apply_async中指定的队列是external_req,但Celery配置定义的队列名称是external_service_req,worker监听的也是这个名称。apply_async的queue参数需要完全匹配配置中的队列名称(而非路由键)。

修改start_distribution中的提交代码:

g.apply_async(
    queue='external_service_req',  # 改为配置中定义的队列名
    expires=distribution.end_date
)

2. 检查任务过期时间是否合法

如果distribution.end_date早于当前时间,组任务会直接被Broker丢弃,不会进入队列。添加日志验证:

from django.utils import timezone

logger.info(f"Current time: {timezone.now()}, Expires at: {distribution.end_date}")
if distribution.end_date < timezone.now():
    logger.warning("Distribution end date is in the past, tasks will not be queued!")

3. 确认待执行任务列表非空

虽然日志打印了client.pk,但可以显式判断是否有任务需要提交:

if not clients_for_send.exists():
    logger.warning("No clients found for distribution, skipping task group submission")
    return

g = group(send_request.s(distribution_id, client.pk) for client in clients_for_send)
g.apply_async(queue='external_service_req', expires=distribution.end_date)

4. 验证Celery配置加载优先级

app.config_from_object会加载Django settings中的Celery配置,如果settings里定义了CELERY_TASK_QUEUES,会覆盖app.conf.task_queues的内容。检查Django settings文件,确保external_service_req队列存在。

5. 直接查看Broker队列状态

用Redis命令确认队列是否存在及任务数量:

# 连接Redis(对应Broker的db 0)
redis-cli -n 0
# 查看所有Celery相关队列
keys celery*
# 查看目标队列的任务数
llen celery_queue_external_service_req

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 18:27:30