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

Celery任务被SSE客户端阻塞,其他.delay调用任务无法触发如何解决?

核心根因
  • 你将SSE消费者函数同时绑定了@worker_ready.connect信号,这个信号会在Celery Worker进程启动完成后直接在Worker的主线程同步执行,@app.task装饰器对信号触发的调用不生效,函数内的SSE长连接死循环会直接占死Worker的调度主线程,导致Worker无法正常拉取、分发队列中其他通过.delay()投递的任务。
  • CELERY_TASK_ALWAYS_EAGER=True配置下所有任务都在调用方线程同步执行,不经过Worker调度,因此不会触发Worker主线程阻塞的问题,表现为运行正常。
可行解决方案

方案1:独立部署SSE消费者(最推荐)

将SSE消费逻辑和Celery Worker完全解耦,单独作为常驻进程运行,不占用Celery的调度、执行资源,完全避免互相影响:

  • 移除SSE消费函数上的@app.task、@worker_ready.connect装饰器
  • 编写独立的启动脚本,直接执行SSE消费逻辑,和Celery Worker分开启动即可

方案2:Celery单独队列运行长驻任务

如果需要保留在Celery体系内管理任务,调整触发逻辑+队列隔离:

  1. 移除@worker_ready.connect装饰器,给SSE任务指定专属队列
@app.task(ignore_result=True, queue="sse_long_running")
def xx_sse_nodes_uptime_info_consumer():
   update_tables...
  1. 单独启动一组Worker进程/线程只消费SSE专属队列,和普通业务队列的Worker完全隔离:
# 启动1个Worker专门处理SSE长驻任务
celery -A your_project worker --concurrency=1 -Q sse_long_running

# 单独启动其他Worker处理普通业务任务
celery -A your_project worker --concurrency=8 -Q default
  1. Worker启动完成后,手动调用一次xx_sse_nodes_uptime_info_consumer.delay()投递长驻任务即可。

方案3:异步非阻塞改造

如果使用线程并发模型,将SSE消费逻辑改为asyncio异步实现,搭配Celery的异步任务支持,避免阻塞工作线程,注意同步实现的长驻任务会持续占用工作线程,就算初始线程数充足,后续如果多启动几个长驻任务依然会把所有工作线程占满,导致普通任务阻塞。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 11:09:03