跳转到内容

Celery

Celery 是一个简单、灵活且绝对可靠的分布式系统,专门用于处理跨进程、跨机器的 大批量异步任务队列 以及 定时周期性任务 的调度。

  • Broker (消息总线/媒介):接收生产者任务并送达消费者的中间人。最经典的载体是 RedisRabbitMQ
  • Beat (定时调度器):类似于 Linux Crontab,负责读取周期任务配置,按时向 Broker 中扔入待执行的消息载荷。
  • Worker (任务消费者):真正在后台物理机上实时监视 Broker 队列、获取并全量执行具体函数逻辑的工作进程。
  • Backend (结果仓库):负责持久化留存任务执行完后的返回状态与 Payload 结果(常见选用 Redis 或 Database)。

(依赖安装:pip install -U celery)


Terminal window
# 启动标准消费者 Worker 进程
celery -A your_project worker -l INFO
# 解决 Windows 环境下原生并发支持的缺陷 (强制使用 eventlet 协程池作为执行单元)
# 需提前 pip install eventlet
celery -A your_project worker -l INFO -P eventlet

在分布式长耗时任务环境中,最怕出现“任务莫名丢失”或“数据库意外断开连接”。以下为实战中必须配置的生命线参数:

3.1. 容错与任务超时重试 (Task Reliability)

Section titled “3.1. 容错与任务超时重试 (Task Reliability)”

默认情况下,Worker 一拉取到任务就会立马向 Broker 回复“确认收到(Ack)”。如果此后 Worker 进程 OOM 崩溃,该任务将永久丢失。

# 核心保命机制:延迟确认。仅当任务函数彻底、成功地执行完毕后,才向 Broker 发送 Ack 信号。
# 配合异常发生时不重投递的特性,可实现崩溃防丢。
CELERY_ACKS_LATE = True
# 设置任务的最长容忍时间,防止死循环任务吃光 Worker 资源
soft_time_limit = 3600 * 24 * 30
# 提升队列中的不可见时间窗口 (Visibility Timeout)
# 对于极端长耗时的任务,若此窗口小于实际运行时间,Broker 会误判 Worker 已死而将任务重新发给其他节点,导致重复执行。
CELERY_BROKER_TRANSPORT_OPTIONS = {'visibility_timeout': 3600 * 24 * 30}

3.2. 长耗时任务引发的数据库 Wait_Timeout 断连

Section titled “3.2. 长耗时任务引发的数据库 Wait_Timeout 断连”

场景:如果一个 Celery 任务在执行计算花费了 3 小时,它从池子里借出的 MySQL 连接在闲置期间超过了 MySQL 全局设置的 wait_timeout(如 8 小时默认,但常被改短),当它计算完准备执行 DB 写入时,连接已被服务端单方面掐断,引发崩溃报错。

根源诊断 (show variables like '%timeout%';):

set global interactive_timeout= 86400;
set global wait_timeout= 86400;

客户端代码防坑解决: 在任务开始执行的钩子,或即将写入数据库前,显式地要求 Django ORM 关闭陈旧连接,强迫其获取新生连接:

from django.db import close_old_connections
@app.task
def long_running_task():
close_old_connections() # 任务启动时清理无效虚假连接
# ... 开始极其漫长的运算 ...

3.3. Worker 远程访问 Redis (Broker) 意外掉线

Section titled “3.3. Worker 远程访问 Redis (Broker) 意外掉线”

在跨内网或公网通过 Socket 直连 Redis 时,极易因网络抖动中断连接:

# 强制调大底层 Socket 握手超时容忍度,并赋予底层自动重连的权利
CELERY_REDIS_SOCKET_CONNECT_TIMEOUT = 5
CELERY_REDIS_RETRY_ON_TIMEOUT = True