Celery队列如何同时拉取多条消息实现独立任务并行执行
Celery短信任务并行执行配置修正方案
核心调整步骤
你当前的并发不生效主要由3个常见错误导致,按优先级调整即可:
1. 修复worker启动命令语法错误
你提供的启动命令存在参数拼接错误,INFO和-Q之间缺少空格,会导致队列指定不生效,worker无法正确拉取目标队列任务,修正后的启动命令如下:
celery worker --app=worker.app --concurrency=5 --hostname=worker1@%h --loglevel=INFO -Q queue1 -Ofair
2. 拆分任务队列,避免调度任务占用发送任务的并发槽位
你当前将定时拉取条目任务get_new_messages和实际的短信发送任务都放在queue1中,调度任务会占用worker的并发资源,导致短信任务无法被及时分配到空闲线程。调整方案如下:
- 新增独立的调度任务队列,修改beat配置:
app.conf.beat_schedule = { "get-message": { "task": "worker.schedule.get_new_messages", "schedule": 10, 'options': {'queue' : 'schedule_queue'} # 单独给调度任务用队列 } }
- 短信发送任务指定到独立的短信队列
sms_queue,更推荐单独启动两组worker分别处理调度任务和短信任务,资源隔离更彻底:
# 处理调度任务的worker,并发数设为1即可 celery worker --app=worker.app --concurrency=1 --hostname=schedule_worker@%h --loglevel=INFO -Q schedule_queue -Ofair # 处理短信发送的worker,并发数按需要调整 celery worker --app=worker.app --concurrency=5 --hostname=sms_worker@%h --loglevel=INFO -Q sms_queue -Ofair
3. 检查get_new_messages的逻辑,确保每个短信任务独立异步提交
常见逻辑错误是:在拉取条目的调度任务中直接串行执行短信发送逻辑,而不是将每个号码的发送任务作为独立的异步任务提交到队列。正确的逻辑示例如下:
@app.task def get_new_messages(): # 从数据库拉取未发送的条目 unsent_sms = SmsModel.objects.filter(status=0).all() for item in unsent_sms: # 每个条目单独提交异步任务,不要在这里直接调用send_sms的同步逻辑 send_sms_task.delay(phone=item.phone, content=item.content) # 标记条目已提交到队列,避免下次轮询重复提交 item.status = 1 item.save() @app.task(queue='sms_queue') def send_sms_task(phone, content): # 实际的短信发送逻辑 send_request(phone, content)
可选优化配置
在Celery配置中新增task_acks_late = True,让worker只有在任务执行完成后才向broker确认消息,配合-Ofair参数可以避免慢任务阻塞队列,更合理地分配任务到空闲线程:
app.conf.update( result_expires=3600, task_track_started=True, worker_prefetch_multiplier = 5, task_acks_late = True # 新增配置 )
完成以上调整后,多个短信发送任务会自动分配到不同的并发线程同时执行,3条短信的总耗时就会降到单条短信的发送耗时水平。
内容的提问来源于stack exchange,提问作者user782400
相关产品推荐
相关产品推荐

