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

