Celery 5.2.6长运行任务被取消及断点续跑方案问询
问题根因
日志里的连接拒绝错误本质是2GB内存配置的Droplet上Redis服务被系统OOM机制杀掉重启,直接导致Celery worker与Broker的连接中断,长任务被终止:
- 本地环境内存充足,不会触发OOM杀进程,因此8-10分钟的长任务可正常执行;2GB基础款Droplet默认未配置swap,Celery worker、beat、Redis、业务进程、系统服务叠加运行时很容易耗尽内存,系统会优先杀掉占用内存较小的用户态进程Redis,触发连接拒绝。
- Celery 5.2.6默认
worker_cancel_long_running_tasks_on_connection_loss = False,但连接中断时若worker进程本身被内存回收、或心跳超时,未完成的长任务既不会继续执行,也不会自动重新入队,只能手动重启触发。
方案一:从根源避免长任务被意外取消
优先解决内存不足导致的进程被杀问题,再通过配置兜底提升容错性:
- 第一步:给Droplet配置2GB swap,从系统层避免OOM杀进程
执行以下命令配置永久生效的swap,配置完成后用free -h可验证swap状态:fallocate -l 2G /swapfile chmod 600 /swapfile mkswap /swapfile swapon /swapfile echo '/swapfile none swap sw 0 0' >> /etc/fstab - 第二步:限制Redis内存占用,避免Redis占满内存挤掉其他进程
编辑Redis配置文件(通常路径为/etc/redis/redis.conf),添加以下配置后重启Redis生效:maxmemory 256mb maxmemory-policy allkeys-lru - 第三步:调整Celery配置适配长任务场景
在Celery配置文件中添加以下参数:# 连接丢失时不直接取消运行中的任务,兼容5.1版本前的默认行为 worker_cancel_long_running_tasks_on_connection_loss = False # 开启broker断连自动重试 broker_connection_retry_on_startup = True broker_connection_max_retries = 100 # 任务执行完成后再确认消费,避免worker异常导致任务丢失 task_acks_late = True # 每次只从队列预取1个任务,避免多任务抢占内存 worker_prefetch_multiplier = 1 # 配置合理的超时阈值,比预期最长任务耗时多1倍以上即可 task_time_limit = 3600 task_soft_time_limit = 3000 - 第四步:调整supervisor配置,实现进程异常自动拉起
给Redis、Celery worker、Celery beat都配置自动重启,长任务单独用低并发worker运行,和短任务做队列隔离:[program:redis] command=/usr/bin/redis-server /etc/redis/redis.conf autostart=true autorestart=true user=redis [program:celery-worker-long] command=/path/to/your/virtualenv/bin/celery -A myproject worker -l info --concurrency=1 -Q long_task directory=/path/to/your/project autostart=true autorestart=true user=www-data [program:celery-beat] command=/path/to/your/virtualenv/bin/celery -A myproject beat -l info directory=/path/to/your/project autostart=true autorestart=true user=www-data
方案二:实现断点续跑,避免任务重跑从头开始
如果极端场景下任务仍被中断,不需要Celery原生支持,通过业务逻辑层的进度持久化即可实现断点续跑,核心逻辑如下:
- 每个长任务绑定唯一标识,任务启动前先从持久化存储(可用现有Redis、业务数据库)查询是否存在历史执行进度
- 将长任务拆分为多个原子执行步骤,每完成一个步骤就将当前进度、中间结果写入持久化存储,标记已完成的步骤
- 任务重新触发时,直接读取上次存储的进度,从未完成的步骤开始执行,跳过已完成的逻辑
最简实现示例:
from celery import shared_task import redis r = redis.Redis(host='localhost', port=6379, db=1) @shared_task(bind=True, queue="long_task") def long_running_task(self, task_id, total_steps=100): # 读取历史进度 progress_key = f"long_task_progress:{task_id}" current_step = r.get(progress_key) current_step = int(current_step) if current_step else 0 # 从断点位置开始执行 for step in range(current_step, total_steps): # 单步业务逻辑 execute_single_step(step) # 每步完成后更新进度 r.set(progress_key, step + 1) # 任务全部完成后清理进度标记 r.delete(progress_key) return "task success"
额外可配置失败信号监听,长任务异常退出时自动重新入队,不需要手动触发:
from celery.signals import task_failure @task_failure.connect def auto_retry_long_task(sender, args, kwargs, **_): if sender.queue == "long_task": sender.apply_async(args=args, kwargs=kwargs, queue="long_task")
内容的提问来源于stack exchange,提问作者ira
相关产品推荐
相关产品推荐

