可靠消息模式与生产实践
投递语义、死信队列、消费幂等、顺序保证与背压控制,把消息系统从"能通"做到"可靠"。
0. 一句话理解
消息系统的可靠性 = 不丢(确认机制)+ 不重(消费幂等)+ 失败有去处(死信队列)+ 快慢有协调(背压)。
1. 不丢消息:三段确认
| 阶段 | 风险 | 对策 |
|---|---|---|
| 生产者 → Broker | 网络闪断,消息没发出 | 发送确认(Kafka acks=all,RabbitMQ publisher confirms) |
| Broker 存储 | 服务器宕机丢数据 | 副本机制(Kafka 多副本、RabbitMQ 镜像/仲裁队列) |
| Broker → 消费者 | 消费者处理失败 | 手动确认:成功才 ack,失败重投 |
# RabbitMQ 生产者确认(pika 示例片段)
channel.confirm_delivery()
if channel.basic_publish(...):
print("Broker 已确认收到")
讲解:
confirm_delivery()开启发布确认:basic_publish返回成功代表 Broker 已持久化,而不是”发出去了”。- 消费者侧”先处理业务、后 ack”;如果先 ack 再处理,处理崩溃就会丢消息。
- Kafka 生产端
acks=all表示副本全部写入成功才确认,配合min.insync.replicas=2使用。
2. 消费幂等:重复无害
# 幂等示例:用唯一键去重(Redis SETNX 或数据库唯一索引)
import redis
r = redis.Redis()
message_id = "order-1001-paid"
ok = r.set(message_id, "1", nx=True, ex=86400)
if ok:
# 第一次处理
process_payment(message_id)
else:
# 重复消息,直接跳过
print("重复消息,忽略")
讲解:
- At-least-once 语义下消息可能重复投递,消费端必须”重复执行也安全”。
- 用消息唯一 ID + 去重存储(Redis
SET NX或数据库唯一索引)是最常用的幂等方案。 - 也可以设计业务天然幂等:如”把余额设置为 100”重复执行结果相同;而”余额 +10”就不幂等。
3. 死信队列:失败消息有去处
业务队列 -> 重试 3 次仍失败 -> 死信队列(DLQ)
-> 人工/定时任务分析补偿
讲解:
- 消息处理失败通常先重试(指数退避),超过次数后放入死信队列,而不是无限重试阻塞主队列。
- 死信队列里的消息由监控告警 + 人工排查,修复后重新放回业务队列。
- RabbitMQ 用
x-dead-letter-exchange声明死信;Kafka 通常用独立的xxx-dlq主题。
4. 顺序保证
必须有序:同一订单的事件(创建 -> 支付 -> 发货)
做法:Kafka 用 key=orderId 保证进同一分区;RabbitMQ 用单队列 + 单消费者
讲解:
- 全局有序成本极高,99% 的场景只需要”业务键内有序”。
- Kafka 一个分区一个消费者实例处理,顺序自然保证;分区数就是并行上限。
- 顺序与吞吐是矛盾:要顺序就别把同一 key 的消息拆到多个分区。
5. 背压与消费能力
生产速率 10 万/秒 > 消费速率 1 万/秒
=> 队列积压 -> 监控告警 -> 扩容消费者 / 优化消费逻辑 / 降级
讲解:
- 消息积压(Lag)是首要监控指标:Kafka 看消费组 Lag,RabbitMQ 看队列 Ready 数量。
- 扩容消费者前先确认瓶颈:数据库慢则加消费者也没用,先优化查询与批量处理。
- 长期积压要考虑降级策略:丢弃非关键消息或合并批量处理,避免”雪崩式追债”。
- Kafka 特有陷阱——重平衡风暴:消费者处理耗时超过
max.poll.interval.ms会被踢出组,触发重平衡、别人也要重投,全组反复抖动。对策:调大该参数、 控制单次 poll 返回量(max.poll.records)、消费逻辑异步化。
6. 本地消息表:把”发消息”变成”写数据库”
第一节的发送确认解决”消息到没到 Broker”,但还有一个更隐蔽的坑: 业务写库成功、发消息失败(或反之),两边永远对不齐。 事务性发件箱(Transactional Outbox)/本地消息表是工程标准答案:
同一本地事务里:
BEGIN
UPDATE orders SET status = 'paid' WHERE id = 1001; -- 业务数据
INSERT INTO outbox (id, topic, payload) VALUES (...); -- 待发消息
COMMIT
后台中继进程(轮询或 CDC):
读取 outbox 未发送行 -> 发到 MQ -> 标记已发送
# 中继进程最小骨架(轮询版)
while True:
rows = db.query("SELECT * FROM outbox WHERE sent = false ORDER BY id LIMIT 100")
for row in rows:
producer.send(row.topic, row.payload).get() # 等确认
db.execute("UPDATE outbox SET sent = true WHERE id = %s", (row.id,))
time.sleep(0.5)
要点:
- 业务写库与”写消息”在同一个本地事务里,要么都成功要么都不成功—— 不依赖任何分布式事务组件。
- 中继失败重发会导致消息重复,回到第 2 节:消费端幂等是整个体系的兜底。
- 生产级实现常用 CDC(Debezium 监听 outbox 表 binlog)替代轮询,延迟更低。
7. 生产检查清单
- 生产者开启发送确认,失败有重试与告警;
- 跨库一致性场景用本地消息表,别赌”先发消息再写库”;
- Broker 开启副本与持久化(Kafka
acks=all+ 副本数 3 +min.insync.replicas=2); - 消费者手动确认 + 重试(指数退避)+ 死信队列;
- 消费逻辑幂等,重复消息无害;
- 监控队列积压、消费 Lag、重试次数、DLQ 深度,全部配告警;
- 消息体带唯一 ID、时间戳与 schema 版本,便于追踪与演进。
8. 动手试试
- 在 RabbitMQ 里搭”业务队列 + 死信队列”:让一条消息处理失败 3 次后进入 DLQ (用 TTL + 死信交换机做延迟重试队列)。
- 用 Redis
SET NX实现一个消费去重器,连续投递同一 ID 两次,确认第二次被忽略。 - 给 Kafka 消费组配一个积压告警:Lag 超过阈值时输出日志
(
kafka-consumer-groups.sh --describe查看 Lag)。 - 用两张表(orders + outbox)模拟本地消息表:故意在中继环节 kill 进程, 验证重启后消息仍会被补发。
9. 一句话记住
可靠 = 确认不丢 + 幂等不重 + DLQ 兜底 + Lag 监控;跨系统一致用本地消息表, 顺序只在业务键内保证,别拿全局有序换吞吐。