消息可靠性设计
消息队列常被当成“发了就一定到”的管道,实际上一段消息要经过生产、存储、消费三个环节,任何一环缺失都可能导致“看似发送成功、实际丢了”。本章按链条拆解每一环该做什么,再补上重试、死信、积压与幂等的处理。
全链路三段
| 环节 | 手段 | 防的问题 |
|---|---|---|
| 生产者 → Broker | 发送结果确认(confirm / acks / SendStatus) | 网络抖动、Broker 拒收导致消息没进去 |
| Broker 存储 | 多副本 + 刷盘策略 + 持久化 | Broker 宕机丢消息 |
| Broker → 消费者 | 手动 ack,处理成功再确认 | 消费失败但消息已被删 |
常见的认知错误是只做了一段就以为可靠:比如生产端开了确认,但消费端用自动 ack,业务抛异常时消息已经确认,照样丢。
生产端:确认 + 本地消息表
先确保“消息真的进了 Broker”:
| 中间件 | 确认方式 |
|---|---|
| Kafka | acks=all 配合 retries,需要防重时可开启生产端幂等(enable.idempotence) |
| RabbitMQ | Publisher Confirm;发布时设 mandatory 捕获不可路由消息 |
| RocketMQ | 同步发送并检查 SendStatus,只有 SEND_OK 才算成功 |
但“确认成功”解决不了业务一致性问题:本地事务提交成功、发消息失败怎么办?用本地消息表把两件事绑在一起。
-- 与业务表同库,保证可以放进同一个本地事务
CREATE TABLE mq_outbox (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
biz_id VARCHAR(64) NOT NULL,
topic VARCHAR(64) NOT NULL,
body TEXT NOT NULL,
status TINYINT NOT NULL DEFAULT 0, -- 0 待投递 1 已投递
retry_count INT NOT NULL DEFAULT 0,
next_time DATETIME NOT NULL,
UNIQUE KEY uk_biz (biz_id)
);
流程:本地事务写业务数据 + mq_outbox → 后台任务(或事务提交后钩子)扫描 status=0 投递 → 投递成功改 status=1;失败按退避更新 next_time。消费端靠 biz_id 幂等,允许重复投递。
Broker 端:持久化与副本
- Kafka:
replication.factor建议 3,min.insync.replicas设为 2(=acks=all才有意义),并关闭“允许非同步副本当 leader”的策略,避免选举出落后副本。 - RocketMQ:同步刷盘更安全、异步刷盘吞吐更高;主从角色配合刷盘策略决定宕机时会不会丢。
- RabbitMQ:队列
durable+ 消息持久化缺一不可;关键业务用仲裁队列(quorum queue)以获得基于 Raft 的副本能力。
“消息持久化了”不等于“不会丢”:如果刷盘异步且 Broker 掉电,仍可能丢最后几条。安全性要求越高,代价越大,按业务分级选择。
消费端:手动 ack
消费端的铁律是业务处理成功后再 ack:
// 以下片段需放进使用 RabbitMQ Java 客户端(或 Spring AMQP)的工程中运行
public void onMessage(Message message, Channel channel) throws IOException {
long tag = message.getMessageProperties().getDeliveryTag();
try {
handle(message); // 业务处理,可能抛异常
channel.basicAck(tag, false); // 处理成功才确认
// 输出:ack 成功,消息从队列移除
} catch (Exception e) {
channel.basicNack(tag, false, true); // 处理失败,requeue=true 重新投递
// 输出:nack,消息稍后重投
}
}
注意事项:
- 不要用自动 ack:消息一发到消费者就确认,进程崩溃时消息直接丢失。
requeue=true要有次数上限,否则一条毒消息(必然失败的坏数据)会无限重投,把队列堵死。- ack 之前就崩,消息会被重投,所以消费逻辑必须幂等。
重试与死信
| 策略 | 做法 | 说明 |
|---|---|---|
| 原地重试 | 同一条消息反复重投 | 只适合临时故障,必须限次 |
| 延迟重试 | 投到重试 topic,按退避时间重新消费 | 缓解拥堵,推荐 |
| 死信队列 | 超过重试上限转入死信 | 避免无限重投,落库待人工处理 |
| 人工补偿 | 后台页面重放死信 | 一定要有,否则死信只能删 |
死信量必须告警:它代表真实业务损失,长期积压说明线上有未修复的问题。
消息积压怎么处理
先判断原因,再动手,否则扩了消费者也没用:
| 现象 | 可能原因 | 处理 |
|---|---|---|
| 某分区堆积,其他正常 | 该 key 热点或单分区消费慢 | 调整分区策略/拆分热点 key |
| 全部堆积且消费 TPS 低 | 消费逻辑有慢查询、外部调用超时 | 优化消费逻辑,加超时与批处理 |
| 消费端日志报错刷屏 | 毒消息反复重试 | 转死信,先让队列继续流动 |
| 突然堆积 | 上游流量激增或下游故障 | 临时扩容 + 降级非核心消费逻辑 |
扩容消费者要注意上限:Kafka 中同一消费者组内,消费者数量超过分区数时多出来的实例会空转(分区是并行度的上限),必要时先扩分区;RocketMQ 也受队列数限制。极端情况下可以用“临时转储”应急:写一个搬运程序把积压消息快速搬到新的临时 topic(按更大分区数),再让大批临时消费者并行处理,事后恢复原拓扑。
# Kafka 查看消费堆积(lag)
kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group order-group
# 输出:各分区的 CURRENT-OFFSET / LOG-END-OFFSET / LAG
幂等消费
重复投递无法避免,只能让消费端不怕重复:
-- 去重表:唯一键 + INSERT IGNORE,天然幂等
CREATE TABLE consume_record (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
msg_key VARCHAR(128) NOT NULL, -- 业务唯一键(订单号等),不用随机的消息 ID
ctime DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_msg_key (msg_key)
);
消费模板:
1. INSERT IGNORE INTO consume_record(msg_key) —— 影响行数 0 说明已消费过,直接 ack
2. 执行业务操作(更新状态、扣减等),尽量用条件更新保证幂等
3. 提交本地事务
4. 成功后 ack
比去重表更简单的一类做法是状态机约束:UPDATE order SET status='PAID' WHERE id=? AND status='UNPAID',重复执行影响行数为 0,天然幂等,也不额外占表空间。
监控指标
| 指标 | 含义 | 告警建议 |
|---|---|---|
| 堆积量 / lag | 未消费消息数 | 持续增长即告警 |
| 消费延迟 | 消息产生到消费完成的时间差 | 超过业务容忍时长告警 |
| 生产失败率 | 发送失败或 nack 比例 | 大于 0 就查 |
| 重试 / 死信数量 | 重试次数与死信堆积 | 死信非零即告警 |
| 消费 TPS | 每秒处理量 | 突降说明消费端异常 |
小结:可靠消息是三层确认拼起来的——生产端确认(必要时用本地消息表兜底)、Broker 多副本与持久化、消费端手动 ack;再配上有限重试、死信队列、积压处理预案和幂等消费。所有环节都做全,才谈得上“不丢不重”。