Celery 的异步任务模型里,任务从调用方发出后不会直接抵达 Worker,而是先进入 Broker 暂存和转发。Redis 能成为 Celery 常用的 Broker 之一,靠的不是专门的消息队列协议,而是 List、BRPOP、Pub/Sub 这类通用数据结构。把 Redis 当作 Broker 时,Celery 会把任务消息序列化后 push 到 List 的尾部,Worker 空闲时通过 BRPOP 阻塞读取,任务确认、超时重投等逻辑由 transport 层补齐。这套方案部署简单、延迟低,但在消息可靠性和持久化上需要额外注意。

Celery Broker 的消息流转机制
Celery 的任务链路可以拆成三个角色:Producer、Broker 和 Worker。调用 delay 或 apply_async 的进程是 Producer,它只负责生成任务消息并交给 Broker;Broker 负责暂存消息,等待 Worker 消费;Worker 消费成功后,再把结果写入 Result Backend。Broker 的选择直接决定了消息能否可靠到达 Worker,也影响重试、延迟任务、定时任务等行为的稳定性。
当使用 Redis 作为 Broker 时,Celery 的 Redis transport 会把任务消息序列化后写入 Redis 的 List。默认情况下,Producer 使用 RPUSH 把消息追加到队列尾部,Worker 通过 BRPOP 阻塞地从队列头部取出消息。这个交互模型非常轻量,Redis 的单线程事件循环可以支撑较高的队列吞吐。任务消息不是简单的函数名加参数,而是一个包含任务名、参数、元数据、headers、时间戳等信息的结构化数据,序列化格式可以是 JSON、pickle、yaml 等。Worker 拿到消息后会反序列化并执行对应任务,执行完成后根据配置决定是否写入 Result Backend。
这里有一份最基本的 Celery 配置,使用 Redis 的 0 号库作为 Broker,1 号库作为结果存储。
from celery import Celery
app = Celery('tasks', broker='redis://127.0.0.1:6379/0')
app.conf.result_backend = 'redis://127.0.0.1:6379/1'
@app.task
def add(x, y):
return x + y
生产环境中不建议只靠一个 Redis 实例同时承担消息队列和结果存储,而是尽量拆分实例或至少使用不同库。因为结果数据通常是临时数据,可以设置过期时间,而 Broker 队列如果被结果数据影响,可能造成内存压力或淘汰。
还要理解确认机制。Celery 默认采用 early ack,也就是 Worker 刚收到消息就向 Broker 确认,这时如果 Worker 进程崩溃,任务仍然会丢失。开启 task_acks_late 后,Worker 会在任务执行完成后再确认,这样能减少任务丢失,但需要配合 visibility_timeout 来防止任务无限占用队列。
用 Redis 作为 Broker 的配置实践
Redis 作为 Broker 的配置集中在 broker_url 和 broker_transport_options 两个地方。前者指定 Redis 地址,支持 redis://、rediss:// 以及 Redis Sentinel 等写法;后者用来控制消息确认、连接超时、重试策略等细节。生产环境建议把序列化器固定为 JSON,避免使用 pickle 带来的安全问题。
app = Celery('tasks', broker='redis://127.0.0.1:6379/0')
app.conf.update(
task_serializer='json',
result_serializer='json',
accept_content=['json'],
task_acks_late=True,
worker_prefetch_multiplier=1,
broker_transport_options={
'visibility_timeout': 7200,
'socket_connect_timeout': 5,
'socket_timeout': 5,
},
broker_connection_retry_on_startup=True,
)
这里 visibility_timeout 是最容易踩坑的参数,它表示 Worker 从 Redis 取走一条消息后,多长时间内没有确认,就把消息重新放回队列。默认值是 3600 秒。如果任务执行时间稳定在 1 小时以上,就必须调大这个值,否则任务会被重复投递。另一类任务是偶发长任务,比如批量同步数据,执行时间可能从几分钟涨到几小时,这种场景更建议把大任务拆小,或者单独给长任务设置更高的超时。
broker_connection_retry_on_startup 也值得显式声明。在较新的 Celery 版本里,如果启动时连不上 Broker,行为会根据该参数变化,若未配置可能直接抛异常,导致 Worker 无法启动。对于使用编排工具拉起 Worker 的场景,建议开启该参数,让 Worker 在 Redis 短暂不可用时能够不断重试,而不是直接退出。
还有 worker_prefetch_multiplier,默认值会让单个 Worker 一次性预取多条消息。对于任务耗时差异较大的队列,预取会让某些 Worker 堆积任务,另一些 Worker 空闲。把该值设置为 1 可以让 Worker 一次只取一条,执行完再取,任务分配更均匀,但要牺牲一点吞吐量。
Redis 与 RabbitMQ 的实际差异
RabbitMQ 是 Celery 官方支持的另一种常见 Broker,它基于 AMQP 协议,具备更强的消息可靠性和路由能力。Redis 的优势在于轻量、部署成本低,很多团队在项目中本来就会部署 Redis 做缓存,直接复用它当 Broker 可以减少一个中间件。但 Redis 的消息模型比 RabbitMQ 简单,可靠性要弱一些,尤其是在持久化和故障恢复方面。
| 对比项 | Redis | RabbitMQ |
|---|---|---|
| 消息协议 | 自定义 List 模型 | AMQP 0-9-1 |
| 消息确认 | 依赖 visibility_timeout 和 ack late | 原生 ACK、NACK、Reject |
| 持久化 | RDB/AOF,可能丢少量数据 | 消息可持久化到磁盘 |
| 路由能力 | 简单队列 | Exchange、Routing Key 灵活路由 |
| 部署复杂度 | 低,通常已有实例 | 高,需要额外维护 |
| 适用场景 | 允许少量丢失或重复的任务 | 不可丢失的关键任务 |
Redis 的持久化并不是为消息队列设计的。RDB 快照会丢失最近的数据,AOF 即使开启 everysec,在极端宕机时也可能丢失 1 秒内的写入。主从切换如果没有配置哨兵或集群,可能会丢消息。RabbitMQ 则可以把队列和消息标记为 durable,配合镜像队列,在节点故障时提供更高的消息可恢复性。
如果业务任务允许重试且具备幂等性,例如发送通知、生成报表、清理缓存等,Redis 作为 Broker 完全够用。如果任务有严格的不可丢失要求,例如支付回调、订单状态流转,RabbitMQ 会更稳妥。另一个重要场景是复杂路由,Celery 的任务队列默认是单一路径,RabbitMQ 可以通过 Exchange 和 Routing Key 实现按主题分发,Redis 在这方面的能力较弱。
消息丢失与重复执行的常见坑
第一个典型的坑是 visibility_timeout 设置过小。假设任务需要执行 2 小时,而 visibility_timeout 保持默认 3600 秒,那么任务执行到 1 小时时,Broker 就会认为这条消息没有被确认,于是重新投递。结果两个 Worker 同时执行同一个任务,产生重复操作。解决办法除了调大超时,还要开启 task_acks_late,并保证任务本身是幂等的。
第二个坑是 Redis 内存淘汰策略。如果 Redis 实例同时承担缓存和任务队列,且 maxmemory-policy 设置为 allkeys-lru,当内存达到上限时,任务消息所在的 List 可能会被淘汰,导致任务凭空消失。对于 Broker 使用的 Redis 实例,建议关闭 LRU 淘汰,或将其设置为 noeviction,让写入失败触发告警,而不是静默丢弃任务。更稳妥的做法是为 Broker 单独部署 Redis,避免其他业务数据影响。
Worker 崩溃也是重复执行的高发来源。希望任务尽量不丢失时,可以同时开启 task_acks_late 和 task_reject_on_worker_lost,这样当 Worker 异常退出后,消息会被重新排队。但要注意,重复执行在这种情况下依然可能发生,所以任务设计必须考虑幂等,例如通过数据库唯一键、Redis 分布式锁或业务侧去重。
app.conf.update(
task_acks_late=True,
task_reject_on_worker_lost=True,
worker_prefetch_multiplier=1,
broker_transport_options={
'visibility_timeout': 10800,
},
)
最后,如果 Redis 同时存储结果,建议为结果设置过期时间,例如 result_expires=3600,避免结果数据无限增长。对于不想保存结果的任务,可以设置 task_ignore_result=True,减少 Redis 的写压力。总体来看,Redis 作为 Celery Broker 是成熟且高效的方案,核心在于理解它的确认边界和持久化限制,再通过幂等设计和监控告警来补齐可靠性。