Celery beat正常入队但运行数天后任务停止执行问题排查
Celery 4.4.0 定时任务长期运行后停止消费故障排查
问题背景
- 技术栈:Python 2 + Celery 4.4.0 搭建秒级定时任务系统,要求任务每秒执行1次
- 部署环境:GCloud Kubernetes集群,1个Redis Pod作为Broker,2个配置完全一致的Celery应用Pod,共用同一份代码与Redis Broker;每个Pod同时启动Beat调度进程与Worker工作进程,可接受任务重复执行,已排除Beat与Worker同Pod部署的架构问题。
故障现象
- 服务连续运行数天后,任务停止实际触发执行,但Beat进程仍每秒正常向队列投递任务
- 重启Pod后服务恢复正常,运行数天后相同故障复现
相关配置
Celery启动命令
celery worker \ --app scheduler \ --without-mingle \ --without-gossip \ --loglevel=DEBUG \ --queues my_queue \ --concurrency=1 \ --max-tasks-per-child=1 \ --beat \ --pool=solo
Celery配置与任务代码
app = Celery(fixups=[]) app.conf.update( CELERYD_HIJACK_ROOT_LOGGER=False, CELERYD_REDIRECT_STDOUTS=False, CELERY_TASK_RESULT_EXPIRES=1200, BROKER_URL='redis://redis.default.svc.cluster.local:6379/0', BROKER_TRANSPORT='redis', CELERY_RESULT_BACKEND='redis://redis.default.svc.cluster.local:6379/0', CELERY_TASK_SERIALIZER='json', CELERY_ACCEPT_CONTENT=['json'], CELERYBEAT_SCHEDULE={ 'my_task': { 'task': 'tasks.my_task', 'schedule': 1.0, # 每秒执行 'options': {'queue': 'my_queue'}, } } ) @task( name='tasks.my_task', soft_time_limit=config.ENRCelery.max_soft_time_limit, time_limit=config.ENRCelery.max_time_limit, bind=True) def my_task(self): print "TRIGGERED"
故障日志
故障出现后Beat进程每秒输出如下调度日志,但无任务实际执行记录:
# 每秒输出 beat: Waking up now. | beat:633 Scheduler: Sending due task my_task (tasks.my_task) | beat:271 tasks.my_task sent. id->97d7837d-3d8f-4c1f-b30e-d2cac0013531
故障定位思路
- 确认Redis队列堆积情况:故障发生时直接连接Redis实例,执行
LLEN my_queue查看队列长度,确认任务是否真的写入队列、是否存在消费完全停止的情况。 - 排查Worker进程阻塞问题:使用的solo池是单线程串行执行模型,若任务执行过程中出现隐式阻塞(比如网络请求无超时、死锁、C扩展调用挂起),会直接导致Worker进程卡死,无法消费新任务,但同进程内的Beat因为调度逻辑和Worker消费逻辑的调度优先级问题,仍会持续输出投递日志。
注意:
solo池为单进程无fork执行模型,配置的--max-tasks-per-child=1参数在该池下不生效,无法通过子进程回收机制规避进程卡死问题 - 排查Redis连接泄漏问题:Celery 4.4.0的Redis传输层存在已知的连接泄漏bug,长时间运行下连接数耗尽后,Worker无法从Broker拉取任务,但Beat投递任务时如果复用未失效的短连接,会出现投递日志正常但消费失败的现象。故障时可执行
CLIENT LIST查看Redis上的连接数,确认是否达到maxclients上限、是否存在大量闲置未释放的Celery连接。 - 排查K8s健康检查配置问题:如果没有配置正确的存活/就绪探针,Worker进程卡死但Pod处于Running状态时,K8s不会自动重启故障实例,故障会持续到人工介入。
解决方案
- 修复Worker阻塞问题
- 给任务内所有外部调用(数据库、缓存、第三方接口)添加显式超时配置,避免无限等待导致进程挂死
- 替换solo池为prefork池,或者将Beat与Worker进程拆分部署,避免单进程内调度逻辑和消费逻辑互相影响
- 修复Redis连接问题
- 在Celery配置中添加Broker连接保活、连接池回收配置:
app.conf.update( # 保留原有配置 BROKER_TRANSPORT_OPTIONS={ 'socket_timeout': 30, 'socket_connect_timeout': 10, 'socket_keepalive': True, 'max_connections': 100 }, BROKER_POOL_LIMIT=10, BROKER_CONNECTION_TIMEOUT=10, ) - 配置Redis的
maxmemory-policy为allkeys-lru,避免内存占满导致的读写异常
- 在Celery配置中添加Broker连接保活、连接池回收配置:
- 添加强制自愈能力
- 给Celery Pod配置存活探针,检测Worker存活状态:
livenessProbe: exec: command: - celery - -A - scheduler - inspect - ping - -d - my_queue initialDelaySeconds: 30 periodSeconds: 10 failureThreshold: 3 - 将
--max-tasks-per-child调整为100(不要设置为1,避免频繁销毁重建进程带来额外开销),在prefork池下定期回收Worker进程,规避长时间运行的内存泄漏、资源泄漏问题
- 给Celery Pod配置存活探针,检测Worker存活状态:
- 临时规避方案:如果暂时无法调整架构,可以给Pod配置定时重启策略,每24小时滚动重启一次Celery实例,避免故障累积出现。
内容的提问来源于stack exchange,提问作者Terry Tú Nguyễn
相关产品推荐
相关产品推荐

