Python 基础体系 · 第 88/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。
Python 消息系统:RabbitMQ、Kafka、消费语义、重试和死信
消息系统解决的不是“把函数放到后台执行”这么单一的问题,而是把生产者和消费者之间的时间、故障和吞吐解耦:
生产者 ──发布──> Broker ──投递──> 消费者 ──处理──> 外部系统
│ │ │
│ │ └── ACK / offset commit
│ └──持久化、路由、分区、重试、保留
└──发布确认
这里至少存在三次需要单独确认的动作:
- 生产者是否把消息交给 Broker;
- Broker 是否把消息交给消费者;
- 消费者是否完成了业务处理。
这三次确认彼此不能互相替代。RabbitMQ 的 publisher confirm 只确认 Broker 接受了发布,consumer acknowledgement 只确认消费者对投递的处理;二者是正交机制。(rabbitmq.com)
一、先建立消息系统的基本模型
1. 消息、生产者、消费者和 Broker
消息是带有数据和元数据的一次传输单元。数据通常是 JSON、Protobuf、Avro 或自定义二进制;元数据可能包含:
message_id:消息唯一标识;event_type:事件类型;occurred_at:事件发生时间;trace_id:链路追踪标识;key:分区或路由依据;attempt:重试次数;schema_version:数据结构版本。
生产者负责创建和发布消息。生产者通常不应该认为“写入本地 TCP 缓冲区”就代表消息已经可靠保存,因为连接可能在 Broker 接收前断开。
消费者负责获取消息并执行处理逻辑。消费者处理消息时可能调用数据库、HTTP 服务、文件系统、支付接口或另一个消息系统。
Broker负责暂存、路由、排序、复制、投递和记录消费进度。RabbitMQ 和 Kafka 都是 Broker,但两者的核心数据模型不同:
- RabbitMQ 更接近“消息投递系统”:消息进入队列,通常被某一个消费者取走并在确认后删除;
- Kafka 更接近“持久化日志”:消息按分区追加保存,消费者通过 offset 表示读取位置,不同消费组可以独立读取同一份数据。
这一区别决定了 ACK、重试、顺序和扩展方式。
二、RabbitMQ 的消息模型:Exchange、Queue、Binding
RabbitMQ 基于 AMQP 0-9-1 时,生产者通常不是直接向 Queue 发布消息,而是先向 Exchange 发布:
Producer
│ publish(exchange, routing_key, body)
▼
Exchange
│ 根据类型和 Binding 路由
├──> Queue A ──> Consumer A
└──> Queue B ──> Consumer B
2.1 Exchange
Exchange 是路由器,不是消息最终消费位置。常见类型包括:
direct:精确匹配 routing key;topic:使用通配符匹配,例如order.*;fanout:广播到所有绑定队列;headers:根据消息头匹配。
Queue 是保存消息并向消费者投递的缓冲区。多个消费者可以共同消费同一个 Queue,此时每条消息通常只会投递给其中一个消费者。
Binding 是 Exchange 到 Queue 的路由关系。生产者只负责提供 Exchange 和 routing key,是否能进入某个 Queue 取决于 Binding。
2.2 RabbitMQ 示例:发布和消费
下面使用 Python 的 pika 客户端演示 RabbitMQ 的基本模型。需要先启动 RabbitMQ,并安装客户端:
python -m pip install pika
创建 rabbitmq_demo.py:
from __future__ import annotations
import json
import os
import uuid
import pika
URL = os.environ.get("RABBITMQ_URL", "amqp://guest:guest@localhost:5672/%2F")
EXCHANGE = "orders"
QUEUE = "order.created.worker"
ROUTING_KEY = "order.created"
def connection() -> pika.BlockingConnection:
parameters = pika.URLParameters(URL)
return pika.BlockingConnection(parameters)
def declare_topology(channel: pika.adapters.blocking_connection.BlockingChannel) -> None:
channel.exchange_declare(
exchange=EXCHANGE,
exchange_type="topic",
durable=True,
)
channel.queue_declare(queue=QUEUE, durable=True)
channel.queue_bind(
exchange=EXCHANGE,
queue=QUEUE,
routing_key=ROUTING_KEY,
)
def publish() -> None:
with connection() as conn:
channel = conn.channel()
declare_topology(channel)
message = {
"message_id": str(uuid.uuid4()),
"order_id": "order-1001",
"amount": 199,
}
channel.confirm_delivery()
channel.basic_publish(
exchange=EXCHANGE,
routing_key=ROUTING_KEY,
body=json.dumps(message).encode("utf-8"),
properties=pika.BasicProperties(
content_type="application/json",
delivery_mode=pika.DeliveryMode.Persistent,
message_id=message["message_id"],
),
mandatory=True,
)
print("published:", message)
def consume() -> None:
with connection() as conn:
channel = conn.channel()
declare_topology(channel)
# 每次最多允许 10 条未确认消息
channel.basic_qos(prefetch_count=10)
def handle(
ch: pika.adapters.blocking_connection.BlockingChannel,
method: pika.spec.Basic.Deliver,
properties: pika.BasicProperties,
body: bytes,
) -> None:
message = json.loads(body)
print("received:", message)
try:
# 这里执行真实业务,例如写数据库
if message["amount"] < 0:
raise ValueError("amount must not be negative")
# 业务成功后 ACK
ch.basic_ack(delivery_tag=method.delivery_tag)
print("acked:", message["message_id"])
except ValueError as exc:
print("permanent failure:", exc)
# 不重新入队,交给死信配置;若没有 DLX,则消息会被丢弃
ch.basic_nack(
delivery_tag=method.delivery_tag,
requeue=False,
)
except Exception as exc:
print("temporary failure:", exc)
# 临时故障可以重新入队,但必须有次数上限,
# 否则会形成高速重试循环。
ch.basic_nack(
delivery_tag=method.delivery_tag,
requeue=True,
)
channel.basic_consume(
queue=QUEUE,
on_message_callback=handle,
auto_ack=False,
)
print("waiting for messages")
channel.start_consuming()
if __name__ == "__main__":
import sys
if len(sys.argv) != 2 or sys.argv[1] not in {"publish", "consume"}:
raise SystemExit("usage: python rabbitmq_demo.py publish|consume")
if sys.argv[1] == "publish":
publish()
else:
consume()
运行:
python rabbitmq_demo.py consume
python rabbitmq_demo.py publish
预期结果类似:
published: {'message_id': '...', 'order_id': 'order-1001', 'amount': 199}
received: {'message_id': '...', 'order_id': 'order-1001', 'amount': 199}
acked: ...
这里有三个关键前置条件:
durable=True只表示 Exchange 和 Queue 的定义可持久化,不代表所有消息自动可靠;DeliveryMode.Persistent表示消息要求持久化,但生产者仍应使用 publisher confirms;auto_ack=False让应用明确决定什么时候完成消费。
RabbitMQ 的手动 ACK 应该在业务处理完成后发送。若消费者在 ACK 前断开连接,未确认消息会被重新入队,因此消费者必须能够处理重复投递。RabbitMQ 也会通过 redelivered 标志提示该投递可能是重新投递,但这个标志不能替代业务幂等。(rabbitmq.com)
三、RabbitMQ 的 ACK、NACK、Reject 和 Prefetch
3.1 ACK 的真实含义
RabbitMQ 的 ACK 不是“我看见了消息”,而是:
当前消费者已经完成了对这次投递的责任,Broker 可以认为这条消息不再需要重新投递。
因此以下代码顺序是危险的:
def unsafe_handle(ch, method, body):
ch.basic_ack(method.delivery_tag)
save_to_database(body)
如果 ACK 成功后进程在 save_to_database() 前崩溃,RabbitMQ 不会再投递消息,而数据库也没有对应记录,结果是消息丢失。
通常应当反过来:
def safe_handle(ch, method, body):
save_to_database(body)
ch.basic_ack(method.delivery_tag)
此时如果数据库已经写入,但 ACK 发送前进程崩溃,消息会再次投递,结果可能是数据库重复写入。因此安全性从“可能丢失”变成了“可能重复”,这正是 至少一次消费 的典型形态。
3.2 basic.nack(requeue=True) 的边界
负确认有两条主要路径:
处理成功 ──> ACK ──> Broker 删除或确认完成
临时失败 ──> NACK(requeue=True) ──> 重新进入队列
永久失败 ──> NACK(requeue=False) ──> DLX 或丢弃
requeue=True 适合临时故障,例如数据库短暂不可用;但它不适合数据格式错误、必填字段缺失、业务规则永久不满足等情况。
错误消息如果每次都立即重新入队,会形成:
取出 ──失败──> 立即入队 ──> 再取出 ──失败──> ...
这会导致:
- 单条坏消息占据消费者;
- CPU 和网络被无效重试消耗;
- 正常消息得不到处理;
- Broker 指标显示吞吐很高,但有效处理量为零。
3.3 Prefetch 是并发窗口,不是线程数
prefetch_count=N 表示一个 Channel 上最多允许大约 N 条未 ACK 投递处于处理中。它控制的是 未确认消息窗口,不是 Python 线程数,也不是 Queue 长度。
若每条消息平均处理时间为 ,消费者并发窗口为 ,理想吞吐上限近似为:
其中:
- 是同时在途的消息数;
- 是单条消息从投递到 ACK 的平均耗时;
- 是每秒可完成的消息数。
例如:
- 单条处理时间 ;
prefetch_count=10。
则理想吞吐约为:
但这是上限,不包含网络、Broker、数据库锁、Python 调度和尾延迟。若把 prefetch 调成 10,000,并不会自动得到一万倍并发;它可能只会让更多消息堆积在消费者进程内存中。
RabbitMQ 文档也明确区分了 prefetch 窗口与消费者处理能力,并指出值为 0 表示不限制未确认消息数量,这在处理不稳定或消息体较大的消费者中可能造成内存压力。(rabbitmq.com)
四、RabbitMQ 的发布可靠性:Persistent 不等于 Confirmed
发布端至少有以下几种状态:
应用调用 publish
│
├── 未到达 Broker
├── Broker 收到但不可路由
├── Broker 路由到 Queue
├── 消息持久化
└── Broker 返回 publisher confirm
Publisher confirm 是 Broker 对生产者的确认。它回答的是:
Broker 是否已经接受并负责这条发布消息?
它不回答:
- 消费者是否处理成功;
- 数据库是否写入成功;
- 下游 HTTP 请求是否成功。
如果消息发布到持久化 Queue,生产者应同时考虑:
channel.confirm_delivery()
channel.basic_publish(
exchange="orders",
routing_key="order.created",
body=b"...",
properties=pika.BasicProperties(
delivery_mode=pika.DeliveryMode.Persistent
),
mandatory=True,
)
mandatory=True 用于让不可路由消息通过返回机制通知生产者;否则消息可能没有进入任何 Queue。Publisher confirm 与 mandatory return 解决的是不同问题:
- confirm:Broker 是否接管了发布;
- return:消息是否存在无法路由的情况。
RabbitMQ 对持久消息的确认时机与 Queue 类型、持久化状态和复制策略有关;应用不应假设 confirm 一定以发布顺序到达,也不应把 confirm 当作业务处理完成信号。(rabbitmq.com)
五、RabbitMQ 重试:不要把“重新入队”当作延迟重试
5.1 直接重新入队的问题
最简单的重试是:
ch.basic_nack(delivery_tag=tag, requeue=True)
它的语义是“现在再试一次”,不是“过 30 秒再试”。如果目标服务需要 30 秒恢复,这种方式会不断快速失败。
更危险的是,RabbitMQ 的 Queue 通常不提供一个通用的“对当前消息睡眠后再次投递”语义。消费者在处理函数中 sleep(30) 会占用一个消费槽位,且 ACK 超时、连接心跳和整体吞吐都可能受到影响。
5.2 延迟队列拓扑
常见方案是使用 TTL Queue 和 Dead Letter Exchange 构造延迟重试:
主队列 orders.main
│ 消费失败,发布到 retry queue
▼
重试队列 orders.retry.30s
│ 消息 TTL 30 秒后过期
│ 过期消息由 DLX 转发
▼
主 Exchange
│
└──> 主队列 orders.main
超过最大次数
│
└──> orders.dlq
示例拓扑声明:
import pika
connection = pika.BlockingConnection(
pika.ConnectionParameters("localhost")
)
channel = connection.channel()
channel.exchange_declare(
exchange="orders.main",
exchange_type="direct",
durable=True,
)
channel.exchange_declare(
exchange="orders.retry",
exchange_type="direct",
durable=True,
)
channel.exchange_declare(
exchange="orders.dlx",
exchange_type="direct",
durable=True,
)
channel.queue_declare(
queue="orders.main.queue",
durable=True,
arguments={
"x-dead-letter-exchange": "orders.dlx",
"x-dead-letter-routing-key": "orders.dead",
},
)
channel.queue_declare(
queue="orders.retry.30s",
durable=True,
arguments={
"x-message-ttl": 30_000,
"x-dead-letter-exchange": "orders.main",
"x-dead-letter-routing-key": "orders.created",
},
)
channel.queue_declare(
queue="orders.dlq",
durable=True,
)
channel.queue_bind(
exchange="orders.main",
queue="orders.main.queue",
routing_key="orders.created",
)
channel.queue_bind(
exchange="orders.dlx",
queue="orders.dlq",
routing_key="orders.dead",
)
connection.close()
失败后不要简单地 requeue=True,而是:
- 给消息增加
attempt; - 把消息发布到
orders.retry.30s; - 确认重试消息发布成功;
- 对原始投递发送
ACK; - TTL 到期后由 RabbitMQ 重新路由回主队列。
伪代码如下:
def handle_failure(ch, method, properties, body):
message = decode(body)
attempt = int(message.get("attempt", 0)) + 1
if attempt > 3:
publish_to_dlq(message, reason="max_attempts_exceeded")
ch.basic_ack(method.delivery_tag)
return
retry_message = {
**message,
"attempt": attempt,
}
publish_to_retry_queue(retry_message, delay="30s")
ch.basic_ack(method.delivery_tag)
这里的 ACK 顺序很重要。若先 ACK 原消息,再发布重试消息,而发布过程失败,消息会丢失;因此生产环境需要对重试发布使用 publisher confirm,只有确认成功后才 ACK 原投递。
另一方面,若重试发布成功后应用在 ACK 前崩溃,原消息可能重新投递并再次发布一条重试消息。这又会造成重复,因此重试消息也必须具备幂等标识,或使用消息 ID 去重。
5.3 TTL 和死信的准确语义
RabbitMQ 的死信并不只表示“处理失败”。消息可能因为以下原因进入 Dead Letter Exchange:
- 消费者以
requeue=False拒绝; - 消息过期;
- Queue 达到长度限制;
- 某些队列类型或策略触发死信。
配置了 DLX 时,Broker 会把符合条件的消息重新发布到指定 Dead Letter Exchange;没有配置时,消息可能被丢弃。(rabbitmq.com)
因此,死信队列不是自动修复队列。它只是把无法继续正常处理的消息保存到一个可观察、可审计、可人工或程序化恢复的位置。
六、死信队列的生产设计
死信消息至少应携带这些信息:
{
"message_id": "msg-123",
"event_type": "order.created",
"payload": {
"order_id": "order-1001"
},
"failed_at": "2026-09-01T10:00:00Z",
"attempt": 3,
"error_type": "ValidationError",
"error_message": "amount must be positive",
"consumer": "billing-worker",
"trace_id": "trace-456"
}
其中 payload 保留原始业务数据,错误元数据用于诊断。不要只记录字符串 "failed",否则死信无法回答:
- 哪条消息失败;
- 失败了几次;
- 由哪个版本的消费者处理;
- 是暂时故障还是永久错误;
- 是否可以安全重放。
6.1 死信处理方式
常见处理流程如下:
DLQ
│
├── 查看原因
├── 修复代码或外部依赖
├── 确认消息是否已部分生效
├── 修改或保留原 message_id
└── 重新发布到主 Exchange
重放之前必须判断业务副作用。例如订单扣款请求已经成功,但消费者在 ACK 前崩溃,此时重放可能再次扣款。正确做法不是盲目重放,而是让扣款接口以 message_id 或业务幂等键实现去重。
RabbitMQ 的死信重新发布本身也不应被误认为带有完整的端到端发布确认;应用若要求死信不丢,应对自己的重放发布过程启用 confirms,并记录重放结果。(rabbitmq.com)
七、Kafka 的消息模型:Topic、Partition、Offset 和 Consumer Group
Kafka 的核心结构是追加日志:
Topic: orders
Partition 0: offset 0 ──> 1 ──> 2 ──> 3 ──> ...
Partition 1: offset 0 ──> 1 ──> 2 ──> ...
Partition 2: offset 0 ──> 1 ──> 2 ──> ...
7.1 Topic 和 Partition
Topic 是逻辑上的消息流名称。
Partition 是 Topic 的有序日志分片。Kafka 只能保证单个 Partition 内的顺序,不能自动保证整个 Topic 的全局顺序。
如果生产者使用同一个 key:
key = order_id
Kafka 通常会把相同 key 的消息映射到同一个 Partition,从而保证同一订单的事件顺序:
order-1001: CREATED → PAID → SHIPPED
但这并不保证不同订单之间的顺序:
order-1001: CREATED
order-1002: PAID
order-1001: PAID
7.2 Offset
每条消息在 Partition 中有一个 offset。消费者有两个容易混淆的位置:
- position:消费者已经拉取到哪里,下一条准备读取什么;
- committed offset:消费者把进度持久化到哪里,重启后从哪里恢复。
Kafka 官方客户端文档明确指出,提交的 offset 应该是“下一条要处理的消息”的 offset,而不是当前消息的 offset。(kafka.apache.org)
假设 Partition 中有:
offset: 10 11 12
message: A B C
消费者处理完 A 后,应该提交:
committed offset = 11
因为 offset 11 是下一条尚未处理的消息。
如果错误地提交 10,重启后会再次读取 A;如果错误地提交 12,重启后会跳过 B。
7.3 Consumer Group
Consumer Group 是一组共同消费 Topic 的消费者。一个 Partition 在同一时刻通常只分配给同一 Group 中的一个消费者:
Topic 有 4 个 Partition
Consumer Group G:
Consumer 1 -> P0, P1
Consumer 2 -> P2
Consumer 3 -> P3
如果消费者数量超过 Partition 数量,多出来的消费者没有 Partition 可消费;因此 Kafka 的并行度上限通常受 Partition 数量约束。
不同 Group 拥有不同消费进度:
orders
├── group=billing 从 offset 100 消费
├── group=analytics 从 offset 80 消费
└── group=warehouse 从 offset 120 消费
这与 RabbitMQ 中“同一个 Queue 的消息通常只交给一个消费者”不同。Kafka 的消息可以按保留策略继续存在,并被多个 Group 独立读取。
八、Kafka 消费语义:自动提交为什么危险
Kafka 常见消费循环可以抽象为:
poll()
│
├── position 前进
├── 业务处理
└── commit offset
关键问题是:业务处理和 offset 提交不是天然的一个原子事务。
8.1 先提交再处理:至多一次
poll message A
commit offset 11
process A
如果 process A 前进程崩溃,重启后从 offset 11 开始,A 不会再读取。
因此:
- A 可能丢失;
- A 不会重复处理。
这接近 至多一次:
8.2 先处理再提交:至少一次
poll message A
process A
commit offset 11
如果处理完成但提交前崩溃:
poll A
process A 成功
进程崩溃
重启后仍从 offset 10 读取 A,于是 A 被再次处理。
此时:
这就是通常更安全的 至少一次。
Kafka 默认允许至少一次语义;禁用生产者重试并在处理前提交 offset,可以构造至多一次语义,但要承担丢消息风险。(kafka.apache.org)
8.3 自动提交的失败窗口
当 enable.auto.commit=true 时,客户端会周期性提交已拉取进度,而不是等待你的业务函数成功。
时间线可能是:
t0: poll() 拉取 A、B、C
t1: 自动提交 offset 13
t2: 处理 A 成功
t3: 处理 B 时进程崩溃
重启后 offset 已经是 13,B 和 C 都不会再次读取。自动提交适合可以容忍丢失、或业务处理本身只是缓存预取的场景;对于订单、账务、库存等业务,通常需要关闭自动提交,改为处理成功后手动提交。
九、Kafka 的重平衡和 poll() 生命周期
Kafka 消费者不仅处理消息,还必须参加 Consumer Group 的协调。消费者长时间不调用 poll(),可能被认为失效并触发重平衡。
两个相关参数分别限制不同问题:
max.poll.interval.ms:两次poll()之间允许的最长处理间隔;max.poll.records:一次poll()最多返回的记录数。
若一次返回 500 条消息,而每条处理 1 秒,那么最坏处理时间约为:
如果 max.poll.interval.ms=300000,也就是 300 秒,消费者可能在处理完这一批之前离开 Group,触发重平衡。
因此不能只看平均吞吐,还要估算:
并满足:
实际处理时间有长尾时,应使用更保守的分位数,而不是平均值。
Python 客户端经常通过后台线程或本地缓冲维持 poll,再把业务处理提交给线程池或进程池。但这样会引入新的问题:提交 offset 时必须确保对应消息真的处理成功,不能因为 poll() 已经返回就提前提交。
十、RabbitMQ 与 Kafka 的消费语义对比
| 维度 | RabbitMQ | Kafka |
|---|---|---|
| 主要模型 | Queue 中的消息投递 | Partition 中的持久化日志 |
| 消费进度 | ACK/NACK | offset commit |
| 消息完成后 | 通常可被删除 | 仍可按保留策略存在 |
| 并发单位 | Queue、Consumer、Channel、prefetch | Partition、Consumer Group |
| 顺序 | 受 Queue、消费者并发和重新入队影响 | 单 Partition 内有序 |
| 重试 | NACK、重新入队、TTL+DLX | 重读 offset、Retry Topic、DLT |
| 广播 | 多个 Queue 绑定同一 Exchange | 多个 Consumer Group |
| 适合场景 | 任务分发、路由、工作队列 | 事件流、日志、回放、流处理 |
不能简单地说“Kafka 比 RabbitMQ 快”或“RabbitMQ 比 Kafka 可靠”。它们优化的抽象不同:
- 如果重点是把任务交给一个 Worker,并在成功后确认删除,RabbitMQ 的 Queue 模型自然;
- 如果重点是保留事件、允许多个系统独立回放和按 Partition 扩展,Kafka 的日志模型自然;
- 如果需要复杂路由、临时队列、请求响应,RabbitMQ 往往更直接;
- 如果需要高吞吐追加写、消费者独立追赶历史数据,Kafka 更合适。
十一、重试策略:按错误类型分类,而不是统一重试
消息失败至少分为三类。
11.1 可恢复错误
例如:
- 数据库连接暂时中断;
- 下游服务返回 502;
- 网络超时;
- 限流响应;
- 依赖服务正在发布。
这类错误可以重试,但应使用退避:
其中:
- :初始延迟;
- :已经失败的次数;
- :最大延迟;
- :随机抖动,避免大量消费者同时重试。
例如 秒、 秒:
第 1 次:约 1 秒
第 2 次:约 2 秒
第 3 次:约 4 秒
第 4 次:约 8 秒
第 5 次:约 16 秒
第 6 次:约 30 秒
11.2 永久错误
例如:
- JSON 无法解析;
- Protobuf schema 版本不兼容;
- 缺少订单号;
- 金额为负数;
- 业务状态不允许执行。
这类消息不应无限重试。重试只能重复得到相同结果,最终应进入 DLQ 或 DLT。
11.3 未知错误
未知异常不能简单归类为“可恢复”。应记录完整堆栈、消息标识、消费者版本和依赖状态,并设置有限重试次数。超过次数后进入死信,否则一个程序 Bug 会让整条消息流不断自旋。
十二、Kafka 的重试和死信:用 Topic 表达状态
Kafka 没有 RabbitMQ 那种由 Queue TTL 和 DLX 直接表达的通用死信拓扑。生产中通常通过多个 Topic 表达重试层级:
orders.main
│ 失败
▼
orders.retry.10s
│ 到期后由转发程序重新发布
▼
orders.main
orders.retry.1m
│
▼
orders.main
orders.dlt
│ 永久失败或超过次数
一种简单的重试消息结构:
{
"message_id": "msg-123",
"original_topic": "orders",
"original_partition": 2,
"original_offset": 1050,
"attempt": 2,
"available_at": "2026-09-01T10:01:30Z",
"error_type": "TimeoutError",
"payload": {
"order_id": "order-1001"
}
}
重试消费者检查 available_at:
from datetime import datetime, timezone
import time
def ready(message: dict) -> bool:
available_at = datetime.fromisoformat(
message["available_at"].replace("Z", "+00:00")
)
return datetime.now(timezone.utc) >= available_at
def retry_loop(consumer, producer) -> None:
for record in consumer:
message = decode(record.value)
if not ready(message):
time.sleep(0.2)
continue
producer.produce(
topic=message["original_topic"],
key=record.key,
value=encode(message["payload"]),
)
producer.flush()
# 只有重新发布成功后,才提交 retry topic 的 offset
consumer.commit()
这个示例用于说明状态转移,生产代码还必须解决两个问题:
produce()成功调用不等于 Broker 已确认,需要检查 delivery callback 或使用客户端提供的确认机制;- 重新发布成功后、提交 retry offset 前崩溃,会导致原消息重复发布,因此目标消费者必须幂等。
更常见的实现是使用专门的重试转发器,或让消费者通过 seek() 重新读取失败位置。但直接停在一个失败 offset 上会阻塞同一 Partition 后续消息:
P0: A(失败) → B(正常) → C(正常)
▲
└── 如果 A 一直重试,B、C 也无法前进
这体现了 Kafka 顺序和重试之间的冲突:
- 保留 Partition 顺序:失败消息阻塞后续消息;
- 追求整体吞吐:把失败消息移到重试 Topic,但可能改变处理顺序。
若业务要求同一订单严格按顺序处理,重试 Topic 设计必须保留同一 key 的路由策略,并避免在旧事件未成功时处理新事件。
十三、消费语义的形式化推导
设:
- :生产者发布;
- :Broker 接收并确认;
- :消费者开始处理;
- :业务副作用成功;
- :消费者确认或提交 offset。
最危险的窗口是:
H ──> S ──崩溃──> C
业务已经成功,但消费进度没有成功保存。恢复后会再次执行 。
另一种危险窗口是:
H ──> C ──崩溃──> S
消费进度已经保存,但业务副作用尚未发生。恢复后不会再执行 。
因此:
C在S前:至多一次,可能丢失;C在S后:至少一次,可能重复。
要得到端到端的“看起来只生效一次”,需要让以下操作原子化:
但 RabbitMQ ACK 或 Kafka offset 通常由 Broker 管理,业务数据库由另一个系统管理,二者不能自动组成同一个本地事务。
13.1 幂等表解决重复
假设消费者要处理消息并写订单账务表,可以建立消费记录表:
CREATE TABLE consumed_messages (
message_id VARCHAR(128) PRIMARY KEY,
consumed_at TIMESTAMP NOT NULL
);
CREATE TABLE account_entries (
entry_id VARCHAR(128) PRIMARY KEY,
order_id VARCHAR(128) NOT NULL,
amount DECIMAL(18, 2) NOT NULL
);
在同一个数据库事务中:
BEGIN;
INSERT INTO consumed_messages(message_id, consumed_at)
VALUES ('msg-123', CURRENT_TIMESTAMP)
ON CONFLICT (message_id) DO NOTHING;
程序检查插入影响行数:
影响 1 行:第一次处理,继续执行副作用
影响 0 行:message_id 已处理,跳过副作用
然后:
INSERT INTO account_entries(entry_id, order_id, amount)
VALUES ('msg-123', 'order-1001', 199.00);
COMMIT;
最后再 ACK RabbitMQ 或提交 Kafka offset。
即使最后一步前崩溃,消息再次到达时,message_id 唯一约束也会阻止重复副作用。
注意:这要求业务副作用也在同一个数据库事务中。若处理逻辑是“调用外部支付接口”,数据库幂等表不能自动回滚外部支付,因此外部接口也必须支持幂等键,或者使用本地 Outbox、状态机和可查询的结果确认。
十四、Outbox:解决数据库提交和消息发布不一致
另一个常见问题是:
写订单数据库成功
发布 Kafka/RabbitMQ 消息失败
或:
发布消息成功
写订单数据库失败
如果数据库和 Broker 不在同一个事务中,直接执行两个动作就存在双写不一致。
Outbox 模式把消息先写入业务数据库:
业务事务
├── 写 orders
└── 写 outbox_events(status='pending')
│
▼
Outbox Relay
│ publisher confirm
▼
RabbitMQ / Kafka
│
▼
更新 outbox_events(status='published')
事务内:
BEGIN;
INSERT INTO orders(order_id, status)
VALUES ('order-1001', 'created');
INSERT INTO outbox_events(
event_id,
event_type,
aggregate_id,
payload,
status
)
VALUES (
'event-1',
'order.created',
'order-1001',
'{"order_id":"order-1001"}',
'pending'
);
COMMIT;
Relay 程序扫描 pending 记录并发布消息。发布成功后再标记为 published。
Relay 在发布成功和更新状态之间崩溃时,会再次发布同一 event_id。因此 Outbox 通常提供的是:
数据库内可靠保存 + 消息至少一次发布
而不是天然 exactly-once。下游仍然需要通过 event_id 幂等。
十五、Exactly-once 到底是什么意思
“Exactly once”必须先说明作用范围。
15.1 投递 exactly-once
如果意思是:
Broker 只把一条消息交给消费者一次,消费者绝不会重复看到它。
那么在存在网络重试、消费者崩溃和 ACK 丢失时,这个目标通常不成立。
15.2 处理 exactly-once
如果意思是:
某个业务副作用在所有故障情况下只执行一次。
那么外部系统必须参与事务或提供幂等协议。
15.3 Kafka 内部 exactly-once
Kafka 支持事务生产和 Kafka Streams 的 exactly-once processing,但它的保证范围主要是 Kafka 输入 offset、Kafka 输出记录以及 Kafka Streams 状态存储之间的原子处理。Kafka 官方设计文档特别提醒,exactly-once 需要阅读适用边界,不能直接推广为所有外部副作用都只发生一次。(kafka.apache.org)
例如:
Kafka input
│
Kafka Streams 事务
├── 提交 input offset
├── 写 Kafka output
└── 更新本地 state store
这三者可以在 Kafka 事务模型中协调。
但如果处理过程还调用:
Kafka input
│
├── 调用支付 API
├── 写 MySQL
└── 写 Kafka output
Kafka 事务并不会自动回滚已经成功的支付 API 或 MySQL 操作。此时仍需要外部系统幂等、事务消息、Outbox 或 Saga 等机制。
十六、Python 中的异步消费者并发
Python 3.14 的 asyncio 提供基于 async/await 的异步 I/O 并发模型,适合网络 I/O 密集型代码;但一个事件循环中的 Python 代码仍然需要避免长时间阻塞,否则其他任务无法运行。(docs.python.org)
一个受限并发的消费者结构可以写成:
import asyncio
from dataclasses import dataclass
@dataclass
class Message:
message_id: str
payload: dict
async def handle(message: Message) -> None:
await asyncio.sleep(0.1) # 模拟数据库或 HTTP I/O
print("processed", message.message_id)
async def worker(
queue: asyncio.Queue[Message],
semaphore: asyncio.Semaphore,
) -> None:
while True:
message = await queue.get()
try:
async with semaphore:
await handle(message)
# 只有 handle 成功后,才确认消息
await ack(message)
except Exception:
await retry_or_dead_letter(message)
finally:
queue.task_done()
async def ack(message: Message) -> None:
print("ack", message.message_id)
async def retry_or_dead_letter(message: Message) -> None:
print("retry or dead-letter", message.message_id)
async def main() -> None:
queue: asyncio.Queue[Message] = asyncio.Queue(maxsize=100)
semaphore = asyncio.Semaphore(20)
workers = [
asyncio.create_task(worker(queue, semaphore))
for _ in range(4)
]
for i in range(100):
await queue.put(
Message(
message_id=f"msg-{i}",
payload={"index": i},
)
)
await queue.join()
for task in workers:
task.cancel()
await asyncio.gather(*workers, return_exceptions=True)
asyncio.run(main())
这里有两层并发限制:
Queue(maxsize=100):限制本地缓冲,避免无限拉取;Semaphore(20):限制同时执行的业务处理数。
如果 RabbitMQ 的 prefetch 是 1,000,而本地只允许 20 个任务处理,其余消息可能已经从 Broker 投递到进程内,但尚未真正执行。若进程此时崩溃,这些未 ACK 消息会重新投递;这通常没问题,但会增加重复和内存占用。
如果使用阻塞式数据库客户端、阻塞式 HTTP 客户端或 CPU 密集计算,不能直接放在事件循环中:
async def bad_handler(message):
result = blocking_database_call(message)
应使用异步客户端,或显式放入线程池/进程池。否则 ACK 延迟、心跳响应和消费吞吐都会受到影响。
十七、与 FastAPI、ASGI 和 Celery 的关系
17.1 FastAPI 请求处理不等于消息消费
FastAPI 运行在 ASGI 服务器上,Web 请求生命周期与消息消费者生命周期不同。不要在每一个请求中临时创建一个 RabbitMQ 或 Kafka Consumer:
HTTP 请求
├── 创建消费者
├── 消费一条消息
├── 关闭消费者
└── 返回响应
这会造成连接创建开销、消费进度不稳定和重复消费。
更合理的关系是:
FastAPI Web 进程
└── 接收请求并发布消息
独立 Worker 进程
└── 长期运行并消费消息
ASGI 应用的 lifespan 可以用于管理应用级资源,但长时间消费任务仍应明确处理取消、关闭和部署副本问题。FastAPI 官方文档和 ASGI 规范分别描述了应用生命周期和异步服务器接口;它们不会自动替你定义 RabbitMQ ACK 或 Kafka offset 语义。
17.2 Celery 与 RabbitMQ/Kafka
Celery 是任务队列框架,Broker 可以使用 RabbitMQ 或其他支持的传输后端。它在底层仍然面对相同问题:
任务被取出
├── Worker 成功执行
├── Worker 失败
├── Worker 在执行后、确认前崩溃
└── Broker 连接中断
因此 Celery 的 ACK、重试、定时和幂等不是独立于消息系统的魔法,而是对上述状态的更高层封装:
- ACK 过早,可能丢任务;
- ACK 过晚,可能重复执行;
- 自动重试需要次数和退避;
- 定时任务仍需防止重复触发;
- 业务函数必须幂等。
如果 Celery 任务调用外部支付、发邮件或修改库存,必须把任务 ID、业务 ID 或幂等键写入业务状态,而不能因为“任务队列保证至少一次”就假设业务只执行一次。
17.3 与 gRPC 的组合
gRPC 常用于消费者调用下游服务。此时消息系统的重试和 gRPC 的重试必须分层设计:
消息重试
└── gRPC 调用
└── gRPC 内部重试
如果 gRPC 客户端一次调用已经重试 3 次,消息消费者又对整条消息重试 5 次,最坏调用次数可能达到:
如果每一层都设置无限重试,系统可能长时间不 ACK,最终触发重复投递或消费者失活。
应区分:
- gRPC
DEADLINE_EXCEEDED:通常是超时,可有限重试; UNAVAILABLE:服务暂时不可用,可退避重试;INVALID_ARGUMENT:请求永久无效,应进入死信;ALREADY_EXISTS:可能表示幂等请求已经成功,需要查询状态,而不是简单重试。
消息处理的总 Deadline 也应覆盖:
数据库读取 + gRPC 调用 + 本地处理 + ACK/offset 提交
若单条消息的最大允许处理时间为 ,则应满足:
否则消费者可能在业务尚未完成时被 Broker 认为失活,造成重平衡或重复投递。
十八、失败表现和诊断方法
18.1 RabbitMQ 消息不断重复
优先检查:
- 是否在业务成功前 ACK;
- 消费者连接是否频繁断开;
- Channel 是否发生协议异常;
requeue=True是否处理了永久错误;- 是否存在多个消费者修改同一业务对象;
prefetch是否远大于实际处理能力。
RabbitMQ 会在连接或 Channel 关闭时重新入队未确认消息,因此重复投递可能是故障恢复的正常结果,而不是 Broker 重复生成了消息。(rabbitmq.com)
18.2 RabbitMQ 消息消失
检查:
- 是否启用了
auto_ack; - Queue 是否 durable;
- 消息是否 persistent;
- 是否启用 publisher confirms;
- Exchange 是否存在正确 Binding;
- 是否使用
mandatory=True处理不可路由消息; - DLX 是否配置错误;
- Queue 是否因 TTL 或长度限制触发死信。
auto_ack 模式下,消息在发送给消费者后就可能被认为已成功交付;如果消费者随后崩溃,消息可能丢失。(rabbitmq.com)
18.3 Kafka 消费延迟不断增加
检查:
lag = log_end_offset - committed_offset
其中:
log_end_offset是 Partition 当前末尾;committed_offset是消费组已提交进度;lag是尚未提交的消息数量。
需要进一步判断:
- 是生产速度超过消费速度;
- 还是某个 Partition 的单条慢消息阻塞;
- 还是消费者频繁重平衡;
- 还是下游依赖变慢;
- 还是失败消息被反复重试。
如果只有一个 Partition lag 很高,增加消费者数量通常没有帮助,因为并行度受 Partition 限制。
18.4 Kafka 消费重复
常见时间线是:
处理消息成功
提交 offset 失败或进程崩溃
重启后从旧 offset 读取
这属于至少一次语义的正常后果。诊断时应通过 topic + partition + offset + message_id 判断:
- 是同一消息重新读取;
- 还是生产者重复发布了两个不同 offset 的相同业务事件;
- 还是消费者错误提交了批次边界。
18.5 Kafka 消费跳过消息
重点检查是否启用了自动提交,或代码是否在业务处理前执行了:
consumer.commit()
还要检查批量处理时是否把一批消息全部视为成功。例如:
for record in records:
process(record)
commit_all(records)
只要中间一条失败而异常被吞掉,后续的整体提交就可能跳过失败消息。
十九、常见误解
误解一:“ACK 表示消息已经被业务系统永久处理”
不一定。RabbitMQ ACK 只表示消费者向 RabbitMQ 确认这次投递;Kafka offset commit 只表示消费进度已保存。业务数据库、支付系统和文件系统是否成功,需要单独保证。
误解二:“持久化消息就不会丢”
RabbitMQ 的 durable Queue、persistent message 和 publisher confirm 分别解决不同层次的问题。缺少发布确认时,生产者无法可靠判断 Broker 是否接管消息。
误解三:“Kafka offset 就是消息 ID”
offset 只在单个 Topic Partition 内有意义。消息迁移到另一个 Topic、Partition 或环境后,offset 不再具有全局身份。业务去重应使用稳定的 message_id 或业务幂等键。
误解四:“重试次数越多越可靠”
重试只能提高暂时性错误的恢复概率。对永久错误,重试次数越多,越会延迟诊断并消耗系统资源。
误解五:“exactly-once 可以让支付只扣一次”
除非支付系统也参与同一事务,或支付接口支持幂等键并能查询历史结果,否则消息 Broker 的 exactly-once 不能保证外部支付只发生一次。
误解六:“提高消费者数量就能提高 Kafka 吞吐”
只有在还有未分配的 Partition 时,增加消费者才可能提高并行度。若 Topic 只有 3 个 Partition,启动 10 个消费者并不能得到 10 路 Partition 并发。
二十、选择和设计的最终落点
可以用下面的判断方式选择模型:
需要任务被一个 Worker 领取并确认完成?
└── 优先考虑 RabbitMQ Queue / Celery
需要事件长期保留、多个系统独立消费和回放?
└── 优先考虑 Kafka Topic
需要复杂路由、广播、临时队列、请求响应?
└── RabbitMQ 的 Exchange 模型更直接
需要按 key 保持顺序并按 Partition 扩展?
└── Kafka 更自然
需要可靠发布?
└── RabbitMQ 使用 publisher confirms
└── Kafka 关注 producer ack、重试、幂等和事务配置
需要可靠消费?
└── 业务成功后 ACK 或提交 offset
└── 接受至少一次带来的重复
└── 通过幂等键消除重复副作用
需要重试?
└── 区分临时错误和永久错误
└── 有限次数、指数退避、随机抖动
└── 超限进入 DLQ/DLT
需要“只生效一次”?
└── 让业务副作用具备幂等性
└── 必要时使用数据库事务、Outbox 或外部幂等协议
消息系统最重要的不是某个客户端 API,而是明确每一个状态的所有权:
谁负责保存消息?
谁负责确认发布?
谁负责确认消费?
失败后消息去哪?
重试由谁调度?
重复副作用如何消除?
死信如何恢复?
RabbitMQ 更偏向确认式投递,Kafka 更偏向可回放日志;两者都不能替代业务幂等。只要把发布确认、消费确认、失败重试和外部副作用分别建模,所谓“可靠消息”就不再是模糊承诺,而可以被拆成可验证的状态转换。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Celery 任务队列:Broker、Worker、ACK、重试、定时和幂等
- 下一篇:Python gRPC:Protobuf、流式调用、Deadline、拦截器和错误模型
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论