消息可靠性设计

消息队列常被当成“发了就一定到”的管道,实际上一段消息要经过生产、存储、消费三个环节,任何一环缺失都可能导致“看似发送成功、实际丢了”。本章按链条拆解每一环该做什么,再补上重试、死信、积压与幂等的处理。

全链路三段

环节手段防的问题
生产者 → Broker发送结果确认(confirm / acks / SendStatus)网络抖动、Broker 拒收导致消息没进去
Broker 存储多副本 + 刷盘策略 + 持久化Broker 宕机丢消息
Broker → 消费者手动 ack,处理成功再确认消费失败但消息已被删

常见的认知错误是只做了一段就以为可靠:比如生产端开了确认,但消费端用自动 ack,业务抛异常时消息已经确认,照样丢。

生产端:确认 + 本地消息表

先确保“消息真的进了 Broker”:

中间件确认方式
Kafkaacks=all 配合 retries,需要防重时可开启生产端幂等(enable.idempotence
RabbitMQPublisher 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 端:持久化与副本

  • Kafkareplication.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;再配上有限重试、死信队列、积压处理预案和幂等消费。所有环节都做全,才谈得上“不丢不重”。

笔记加载中…