系统架构 · 事务消息与最终一致性

跨服务的一致性拼图: 能不用分布式事务就不用 — 本地事务保原子, 消息保可达, 幂等保不重, 补偿保回退, 对账保最终兜底

事务消息 · Outbox 时序 × Saga 补偿链 × 投递语义 — 双簿记消灭双写分叉 本地事务保原子 · 消息保可达 · 消费端幂等保不重 ① Outbox 发件箱 — 业务表与消息同一事务落盘, relay 异步搬运, 消费端去重 ② 轮询/CDC ③ publish ④ 至少一次投递 订单服务 · 本地事务 orders 订单表 outbox 发件箱表 (biz_id) INSERT 两表 → 同一个 COMMIT 崩在发送前? 事件还在表里 Relay / Debezium 轮询 sent=0 或监听 binlog 发送成功才标记已投 Kafka / RocketMQ at-least-once 投递 重投是常态, 不丢 消费者服务 inbox 去重表 (msg_id) INSERT inbox 撞唯一键? → 已消费过, 直接 ACK ✗ 双写分叉: 先 COMMIT 订单、再发 MQ — 第二步一崩, 库存永远收不到"已下单", 少发货事故就是这么来的 ✓ Outbox 把"发消息"变成"多写一行", 用本地事务的原子性换投递的可靠性; 消费端用 inbox 唯一键挡重复 ② Saga 补偿链 — 正向 T1→T2→T3, 失败反向 C2→C1; 补偿是业务反操作, 不是 rollback T1 T2 T3 ① 创建订单 状态 INIT ② 扣款 冻结/划扣 50 元 ③ 扣库存 deduct(sku, 1) 订单完成 终态 CONFIRMED 库存不足, T3 失败 补偿 补偿 漏补 ③ 失败: 库存不足 Saga 触发反向补偿 C2 取消扣款 退款/解冻 — 必须幂等 C1 关闭订单 状态改 CLOSED T+1 对账发现分叉 库存已扣, 订单却是 CLOSED ✗ 事故现场: 支付渠道超时, C2 退款调用失败且无人重试 — 资金冻结 48 小时, 客诉爆发才被发现 ✓ 正解: 补偿动作可重试 + 幂等; 补偿也失败 → 落人工干预表 + 告警, 让人兜底而不是让数据悬挂 事故现场: org.apache.kafka.clients.consumer.CommitFailedException: Commit cannot be completed since the group has already rebalanced… — 堆积时盲目重启消费者 ③ 投递语义对比 — MQ 只承诺至少一次; "恰好一次"永远是你自己在消费端做出来的 At-most-once 至多一次 发完不等确认, 丢了就丢了 绝不重复, 但可能丢 适用: 指标打点 / 心跳上报 关掉重试与 ACK 即近似此语义 At-least-once 至少一次 不丢, 但必然可能重 Kafka / RocketMQ / RabbitMQ 默认 重投场景: 消费超时 / rebalance ⇒ 消费端必须幂等, 这是底线 Exactly-once 恰好一次 Kafka EOS: 幂等 producer + 事务 isolation.level=read_committed 只覆盖 Kafka 流内部, 不覆盖出口 写外部系统? 还是要幂等键 读法: emerald 正向链 / rose 事故与补偿 · Outbox 解决"发不发得出", 幂等解决"重不重", Saga 解决"回不回得去" 对账是最终一致性的最后一道闸: 任何跨系统链路都要有 T+1 对账兜底, 差异进人工干预闭环

机制视角 — 一致性谱系

  • • 单库 ACID 免费送, 跨了服务就只剩 BASE
  • • 2PC 强一致但同步阻塞, 高并发场景基本不用
  • • TCC 用业务预留换性能, Saga 用补偿拆长事务
  • • Outbox/Inbox 双簿记: 发送原子化 + 消费去重化

行为视角 — 重投是常态

  • • MQ 只承诺 at-least-once, 重复要靠消费端挡
  • • exactly-once 只在 Kafka 流内部成立, 出口仍要幂等
  • • 有 Saga 必须有补偿, 补偿失败必须升级人工
  • • DLQ 是缓冲带不是垃圾桶, 堆积无人管等于丢消息

生产价值 — 对账兜底

  • • 任何跨系统链路都要 T+1 对账做最后防线
  • • 消息必带业务 ID, 对账才有 join 键
  • • 积压预案: 扩容到分区数 + 跳过 + 离线补数
  • • 幂等键选业务自然键, 唯一索引是最后闸门

💡 一句话理解

把跨服务一致性想成两个部门的报销流程: 2PC 是"两边同时签字才生效", 任何一边卡住全体干等; TCC 是"先占额度, 再正式划扣, 不行就释放"; Saga 是"先各办各的, 出了岔子按流程冲销"; Outbox 则是"把要传的话写进同一个小本本, 专人负责抄送"。

机制本质: 分布式事务的终极答案不是"更强的 2PC", 而是把长事务拆短 — 本地事务保原子, 消息保可达, 幂等保不重, 补偿保回退, 对账保最终兜底。能最终一致就别强一致: 一致性是有价格的, 且价格随并发指数上涨。

🧠 必知必会 必考 & 必会

ACID 与 BASE
单库事务靠 ACID (原子/一致/隔离/持久); 跨服务后锁不住也提交不了, 只能退到 BASE (基本可用 + 软状态 + 最终一致)。BASE 不是放弃一致性, 是把"一致"延迟到未来某刻。
BEGIN;
  UPDATE account SET bal = bal - 50 WHERE uid = 1;
  UPDATE account SET bal = bal + 50 WHERE uid = 2;  -- 关键: 同库才能同事务
COMMIT;
-- → 跨了库/服务, 这两条就没人替你保证同生共死
2PC 两阶段提交
协调者先问所有参与者"能否提交" (prepare, 各自锁资源), 全票通过再发 commit。两大死穴: prepare 后资源被锁住干等 (同步阻塞), 协调者单点, 崩在中间全体卡死。
阶段1: coordinator → 全体: prepare?  参与者锁行、写 redo, 回 yes
阶段2: 全票 yes → commit; 任一 no → rollback
-- 关键: 协调者崩在两阶段之间, 参与者抱着锁干等
-- → 协调者单点 + 同步阻塞, 高并发场景基本弃用
3PC (概念层)
在 2PC 中间插入 CanCommit, 并给参与者加超时: 协调者失联时参与者超时后自行推进, 缓解阻塞。但网络分区下仍可能两边各自提交, 实践很少用, 更多是理解 2PC 缺陷的台阶。
2PC: prepare ── commit            (参与者无超时, 卡死等协调者)
3PC: canCommit ── preCommit ── doCommit (超时可自行推进)
-- 关键: 用超时换掉阻塞, 但分区下仍可能双提交
TCC 资源预留
业务层的两阶段: Try 预留资源 (冻结 50 元不真扣), Confirm 划扣, Cancel 解冻。三段都要幂等, 还要处理空回滚 (没收到 Try 先收到 Cancel) 和悬挂 (Cancel 之后 Try 才到)。
@TwoPhaseBusinessAction(name = "deduct", commitMethod = "confirm",
                        rollbackMethod = "cancel")   // Seata 1.5+
boolean tryDeduct(BusinessActionContext ctx, Long uid, BigDecimal amt);
boolean confirm(BusinessActionContext ctx);  // 划扣冻结额, 必须幂等
boolean cancel(BusinessActionContext ctx);   // 关键: 空回滚也要返回成功
// → Try 冻结 50, Confirm 扣 50, Cancel 退冻结
Saga 补偿链
把长事务拆成一串本地事务 T1→T2→T3, 每步配一个业务补偿 C1/C2, 失败时反向补偿。补偿是"业务反操作" (退款、关单), 不是数据库 rollback — 已提交的事务回不去了。
正向: T1 创建订单 → T2 扣款 → T3 扣库存
T3 失败 ⇒ C2 退款 → C1 关闭订单
-- 关键: 补偿必须可重试且幂等, 失败转人工干预
-- → 最终回滚到 T1 之前的业务等价状态
Outbox 发件箱
"写业务 + 发消息"横跨 DB 和 MQ 两个系统, 永远无法原子。Outbox 把消息变成一行数据, 和业务表同一事务落盘; relay 轮询或 CDC 把它搬给 MQ。事件不丢, 只是"可能晚到、可能重"。
BEGIN;
  INSERT INTO orders(...);
  INSERT INTO outbox(event_id, payload) VALUES('evt_1', '{...}');
COMMIT;  -- 关键: 消息和业务同事务, 原子性免费送
-- relay: SELECT * FROM outbox WHERE sent = 0 ORDER BY id LIMIT 100
Inbox 去重收件箱
MQ 的 at-least-once 决定重复投递必然发生; 消费端用 inbox 表的唯一键当闸门: 先 INSERT 消息 ID, 撞键说明消费过, 直接 ACK 跳过。
INSERT INTO inbox(msg_id, consumed_at) VALUES('evt_1', NOW());
-- → Duplicate entry 'evt_1' for key 'PRIMARY'
-- 关键: 唯一键冲突 = 已消费, 捕获后直接 ACK 跳过
At-least-once 投递
broker 没收到 ACK 就重投, ACK 在网络中丢失也重投 — 重复是常态不是异常。所有消费端代码都要按"这条消息我可能见过"来写。
生产者 → broker: 未收到 ack → 重发 → 消费者收到 2 次
消费者 → broker: 处理完未 ack 就崩 → rebalance → 再投 1 次
# 关键: at-least-once 是"不丢"的代价, 重投无法根除只能去重
Exactly-once (EOS)
Kafka 的 EOS = 幂等 producer (PID + 序列号, 单分区会话内去重) + 事务 (跨分区原子写 + offset 提交绑定) + 消费端 read_committed。只覆盖"Kafka 进 Kafka"闭环, 写外部系统仍要幂等。
enable.idempotence=true        # PID + seq: 单分区内不重不乱
transactional.id=svc-1          # 跨分区原子 + 崩溃后 fencing 旧实例
# 关键: consumer 端 isolation.level=read_committed 才配套生效
幂等键 Idempotency
让"执行 N 次 = 执行 1 次"的业务唯一标识: 订单号、支付流水号。落库靠唯一索引兜底, 谁都能重试, 结果唯一。
ALTER TABLE payment ADD UNIQUE KEY uk_order_no (order_no);
-- 重复请求 INSERT → 1062 Duplicate entry 'PO20260901'
-- 关键: 幂等键用业务自然键, 自增 ID 每次都不同挡不住
Deduplication 去重
与幂等一体两面: 幂等是"重复执行无害", 去重是"重复的直接丢弃"。去重记录要有 TTL/清理策略, 不然表无限膨胀; 高并发可先用 Redis SETNX + 过期挡一道, DB 唯一键终审。
-- 先 Redis: SETNX msg:m_88 EX 86400 挡洪峰, DB 唯一键终审
INSERT INTO dedup(msg_id, expire_at) VALUES('m_88', NOW() + INTERVAL 1 DAY);
-- 关键: 去重记录必须带过期清理, 否则 90 天后慢查询拖垮库
Retry Queue 重试队列
消费失败不该原地无限重试 — 进梯度重试队列 (5s/30s/5m/30m), 每级 TTL 到点再投。RabbitMQ 用 TTL+DLX 实现, Kafka 用延迟 topic + 定时轮询。
失败 → retry.5s → retry.30s → retry.5m → DLQ
# 关键: 每级都设 TTL + 死信路由, 到点自动晋级, 不死循环
# → 第 4 级还失败就进 DLQ, 等人工/定时修复重放
DLQ 死信队列
重试耗尽、被 reject、TTL 过期消息的去处。DLQ 不是垃圾桶: 要监控深度、要有人工或自动的修复重放闭环, 否则等于"体面地丢消息"。
QueueBuilder.durable("order.work")
    .withArgument("x-dead-letter-exchange", "order.dlx") // 死信去向
    .withArgument("x-message-ttl", 60000)                // 60s 后进 DLX
    .build();  // 关键: DLQ 深度告警 + 修复重放脚本先备好
Ordering 顺序保证
Kafka 只保证分区内有序: 同一业务 key (订单号) 进同一分区。坑在重试: max.in.flight > 1 且未开幂等时, 第 1 条失败重试会让第 2 条先到 — 开 enable.idempotence 才能在 ≤5 in-flight 下保序。
# 同一订单的消息带同 key → hash 落同一分区, 分区内 FIFO
producer.send(new ProducerRecord("orders", orderId.getBytes(), payload));
enable.idempotence=true   # 关键: in-flight≤5 时失败重试也不乱序

🏭 生产实战 real world

场景 1 · 每月几十单少发货清零 — Outbox 表 + 轮询 relay 上线

之前"先写订单再发 MQ", 发送失败就丢事件。改成同事务写 outbox, 后台 relay 扫表投递, 投递失败只标记不删除。

CREATE TABLE outbox_event (
  id         BIGINT AUTO_INCREMENT PRIMARY KEY,
  event_id   VARCHAR(64)  NOT NULL,           -- 幂等键: 全局唯一
  agg_type   VARCHAR(32)  NOT NULL,           -- 'ORDER'
  agg_id     VARCHAR(64)  NOT NULL,           -- 业务键: 订单号
  payload    JSON         NOT NULL,
  sent       TINYINT      NOT NULL DEFAULT 0,
  created_at DATETIME     NOT NULL,
  UNIQUE KEY uk_event (event_id),
  KEY idx_relay (sent, id)
);
-- 业务代码: 同事务 INSERT orders + INSERT outbox_event
-- relay (MySQL 8.0+): SELECT ... WHERE sent=0 ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED
-- 发 Kafka 成功后: UPDATE outbox_event SET sent=1 WHERE id=?
-- 关键: 投递失败只标记不删除, 下轮继续 — 事件永不丢

场景 2 · 秒杀下单扣库存 — RocketMQ 事务消息: 半消息 + 本地事务 + 回查

要求"订单失败库存一定不扣"。先发半消息, 执行本地事务, 成功 commit; broker 收不到确认就回查本地事务状态, 彻底消灭两步不一致。

// rocketmq-client 4.9.x: 下单扣库存事务消息
TransactionMQProducer producer = new TransactionMQProducer("order-tx-group");
producer.setTransactionListener(new TransactionListener() {
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        boolean ok = orderService.createWithDeduct((Order) arg); // 本地事务
        return ok ? LocalTransactionState.COMMIT_MESSAGE
                  : LocalTransactionState.ROLLBACK_MESSAGE;
    }
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        // broker 收不到确认就回查: 按订单号反查本地事务事实
        return orderMapper.exists(msg.getKeys())
            ? LocalTransactionState.COMMIT_MESSAGE
            : LocalTransactionState.UNKNOW; // 查不清先 UNKNOW, 下轮回查
    }
});
producer.sendMessageInTransaction(msg, order); // 关键: 半消息先行, 回查兜底

场景 3 · 每分钟 GMV 统计不再重复计数 — Kafka 端到端 exactly-once 流处理

消费 offset 提交和产出消息曾是两步, 崩在中间就算重。用 Kafka 事务把"消费位移 + 产出消息"绑成一个原子操作。

props.put("transactional.id", "gmv-agg-1");  // 实例唯一, 崩溃后 fencing 旧实例
props.put("enable.idempotence", true);
producer.initTransactions();

var records = consumer.poll(Duration.ofSeconds(1));  // read_committed 消费
producer.beginTransaction();
produceAggregates(producer, records);                // 产出聚合结果
// poll 循环里累积 pendingOffsets: Map<TopicPartition, OffsetAndMetadata>
producer.sendOffsetsToTransaction(pendingOffsets, consumer.groupMetadata());
producer.commitTransaction(); // 关键: 崩在任意点, 重放不重算不重写
// 消费端必须 isolation.level=read_committed, 才看不到未提交结果

场景 4 · 重投不再引发重复发货 — 全公司统一的幂等消费模板

所有消费者按同一模板写: 去重表 INSERT 与业务处理同一事务, 唯一键冲突即跳过。从此 at-least-once 的"重"被关进笼子。

-- 消费入口 (任一 MQ 通用):
INSERT INTO inbox_consume (msg_id, topic, consumed_at)
VALUES ('evt_20260901_88', 'order.paid', NOW());
-- 冲突(1062) → 已消费过 → 直接 ACK return
-- 未冲突 → 继续执行业务 (扣库存/发券), 与上面 INSERT 同一 DB 事务
COMMIT;  -- 关键: 去重记录和业务结果同事务, 崩了一起回滚, 重投再来

场景 5 · 解析失败的消息曾无限重入队 — RabbitMQ 死信队列 + 修复重放闭环

毒消息曾把消费者 CPU 打满。配置 DLX 收容死信, 手动 ACK + 有限重试, 深度告警后由脚本修复重放。

spring:
  rabbitmq:
    publisher-confirm-type: correlated   # 发送端: broker 确认落盘
    publisher-returns: true
    listener:
      simple:
        acknowledge-mode: manual          # 手动 ack, 处理完才确认
        default-requeue-rejected: false   # 关键: 异常消息进死信, 不回队
        retry:
          max-attempts: 3
          initial-interval: 2000

场景 6 · 跨账户转账资金链路 — Seata TCC: 冻结/划扣/解冻, 空回滚与悬挂都堵上

Try 校验余额并冻结、记 xid 流水; Cancel 按 xid 幂等退冻结, 没有流水就是空回滚, 直接返回成功。

@Transactional
public boolean tryDeduct(BusinessActionContext ctx, Long uid, BigDecimal amt) {
    if (frozen.exists(ctx.getXid())) { return true; }  // 幂等: Try 重入直接成功
    int n = jdbc.update("UPDATE account SET bal = bal - ?, frozen = frozen + ? "
        + "WHERE uid = ? AND bal >= ?", amt, amt, uid, amt); // 余额不足即 Try 失败
    if (n == 0) { return false; }
    frozen.insert(ctx.getXid(), uid, amt);  // 记冻结流水, 防悬挂依据
    return true;
}
@Transactional
public boolean cancel(BusinessActionContext ctx) {
    // 空回滚: Try 没成功过 → 不加钱直接返回成功 (同款 xid 判断防悬挂)
    return frozen.refundByXid(ctx.getXid()); // 关键: 按 xid 退冻结, 幂等可重试
}

场景 7 · "已发货"把订单改回"待支付" — 同 key 同分区修好事件乱序

没设 partition key 时轮询落分区, 同一订单的事件无序到达。改 producer 按 orderId 做 key, 同单事件进同一分区 FIFO。

// 乱序根源: 不设 key → 轮询分区, "已支付"落 P0, "已发货"落 P2
ProducerRecord<String, String> rec =
    new ProducerRecord<>("order-events", orderId, toJson(event));
// orderId 做 key: 同一订单全部事件进同一分区, 分区内 FIFO
// 关键: 消费端再按 (orderId, version) 比较版本, 只前滚不回滚

场景 8 · 积压 400 万条的深夜应急 — 看积压、扩容、跳过三步走

上游 bug 灌入脏数据, 消费者全卡在重试, lag 四小时涨到 400 万。先看积压, 再扩容到分区上限, 脏数据快速跳过转死信留证。

# 1. 看积压: consumer group 的 lag
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
  --describe --group order-consumer   # → LAG 4,213,376
# 2. 扩容: 消费实例 ≤ 分区数 (12 分区最多 12 实例, 再多空转)
kubectl scale deploy order-consumer --replicas=12
# 3. 脏数据快速跳过: 解析失败不重试, 转死信 topic 留证
# 关键: 先止血(跳过)再补数(离线回放), 别让正常业务陪葬

场景 9 · 十年老服务"写库后发 MQ"不敢动 — Debezium 订阅 binlog, 变更即事件

不动业务代码: Debezium 2.x 订阅 binlog, 已提交的行变更自动转事件进 Kafka, 双写分叉从根上消除。

{
  "name": "order-cdc",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "mysql-primary",
    "database.include.list": "shop",
    "table.include.list": "shop.orders",
    "topic.prefix": "shop-cdc",
    "snapshot.mode": "schema_only",
    "tombstones.on.delete": "false"
  }
}
# 关键: CDC 把"已提交的行变更"变成事件流, 天然无双写分叉

场景 10 · 头部商户挤爆一个分区 — 顺序消息的热点打散: 片内有序

按商户 ID 分区保序, 结果头部商户独占 40% 流量把单分区打满。打散: 商户内按订单再分片, 全局有序降级为片内有序。

// 错: key = merchantId    // 头部商户占该分区 40% 流量, 消费者单核打满
// 对: 分片 = orderNo.hashCode() & 7   // 商户内按订单拆 8 个逻辑片
//     key  = merchantId + "#" + 分片
// 全局有序降级为片内有序: 同一订单仍严格有序, 不同订单可并行消费
// 监控: 每分区消息数方差 < 20% 才算打散成功
// 关键: 片数取 2 的幂, 位运算代替取模, 分片更均匀

⚠️ 编码注意与常见坑 pitfalls

坑 1 · 库存少扣了, MQ 里压根没这条消息 — 症状: 下游凭空缺事件. 原因: 先写库再发 MQ 的双写分叉, 写库成功、发消息失败, 没有任何机制补偿. 正解: Outbox/事务消息, 让消息与业务同事务。
-- 错: COMMIT 订单 → send(mq) 失败 → 库存永远不知道这一单
-- 对: 同事务 INSERT outbox → relay 异步投递, 崩了也不丢事件
--     消费端再配 inbox 去重, 重投也不重做
坑 2 · 开了 Kafka 事务, 消费端照样重复扣款 — 症状: EOS 开了钱还是重扣. 原因: Kafka EOS 只保证"Kafka 流内部"不丢不重, 从 Kafka 写到 DB/HTTP 的那段不在保护范围. 正解: 出口必须自己做幂等。
# 错: enable.idempotence=true 就以为万无一失 → 消费写 DB 仍重复
# 对: DB 侧 inbox 去重表 + 唯一键, 消费端幂等自己做
#     EOS 只管流内部, 出口是另一个系统的事
坑 3 · 补偿走到一半发现没有"取消扣款"这一步 — 症状: Saga 卡在中间态. 原因: 设计正向动作时没给每步定义反操作, 失败后无补偿可走. 正解: 每个本地事务 T 强制配一个可幂等重试的 C。
# 错: T1 建单 → T2 扣款 → T3 扣库存(失败) → 无 C2 可走, 卡死
# 对: T 配 C 成对设计: C1 关单 / C2 退款 / C3 回补库存
#     上线前评审表: 每个 T 的 C 是什么, 写进设计文档
坑 4 · 退款调了三天没成功, 没人知道 — 症状: 资金悬挂到客诉爆发. 原因: 补偿也依赖第三方渠道, 渠道挂了补偿同样失败, 只 log 不升级. 正解: 补偿失败落人工干预表 + 告警 + 看板。
-- 错: if !refund() { log.error("cancel fail") }  -- 沉没在日志里
-- 对: INSERT INTO manual_task(xid, action, retry_cnt)
--     -- 关键: 补偿失败必须升级为人 + 定时重推
坑 5 · 死信队列里躺着 30 万条消息三个月 — 症状: DLQ 只进不出. 原因: DLX 配了就不管, 没有深度告警和重放工具, 死信等于体面丢单. 正解: DLQ 深度告警 + 修复重放例行化。
# 错: 配了 DLX 就不管 → 30w 条死了 3 个月才发现
# 对: 深度>100 告警; 修复后按原 exchange 重新 publish 回业务队列
#     每周固定重放窗口, 重放前先修数据再回放
坑 6 · "已支付"后到, 订单状态被回写成"待支付" — 症状: 状态机倒转. 原因: 没设 partition key, 轮询落分区, 同一订单事件无序. 正解: 业务 key (orderId) 做 partition key, 同 key 同分区 FIFO。
// 错: new ProducerRecord("orders", payload)          // 无 key 轮询
// 对: new ProducerRecord("orders", orderId, payload) // 关键: 同 key 同分区
坑 7 · 发一张 8MB 的 Excel 当消息, broker 拒收还拖垮同分区 — 症状: 大消息生产失败或复制变慢. 原因: Kafka 默认 message.max.bytes 约 1MB, 盲调大后复制带宽与消费延迟齐涨. 正解: 传引用, 消息只带元数据, 文件走对象存储。
# 错: message.max.bytes=10485760  # 调到 10MB, 复制带宽翻 10 倍
# 对: 消息体只放 {"fileUrl":"oss://..."}, 文件走对象存储
#     消费端按 URL 拉取, 失败可重试且不占 broker
坑 8 · 消息显示已回滚, 订单其实成功了 — 症状: 事务消息状态与库不一致. 原因: 回查接口实现成读内存/永远返回 UNKNOW, 不查"本地事务的最终事实". 正解: 回查按 msgKey 查库: 有单 COMMIT, 查不清 UNKNOW 两次后再 ROLLBACK。
// 错: 永远 return UNKNOW → 半消息悬挂, 3 次回查后丢弃
// 对: 查库有单→COMMIT; 无单→回查计数<3 return UNKNOW, 之后 ROLLBACK
坑 9 · 重复消息每条都当新的处理 — 症状: 去重完全不生效. 原因: 幂等键用数据库自增 ID, 每次生成不同, 重复消息永远撞不上唯一键. 正解: 幂等键用业务自然键 (订单号/流水号), 唯一索引兜底。
-- 错: id BIGINT AUTO_INCREMENT 做唯一键 → 重投 100 次插 100 行
-- 对: UNIQUE KEY uk_biz (order_no)  -- 关键: 业务键才拦得住重复
坑 10 · 一条毒消息在重试队列里循环了 6 个小时 — 症状: 消费者 CPU 打满无产出. 原因: 失败回原队列且无 TTL/次数上限, 毒消息永动. 正解: 梯度重试 + 最大次数/TTL, 超限进 DLQ。
# 错: 失败 requeue 回原队列 → 毒消息死循环, CPU 100%
# 对: retry 队列 TTL 5m → 3 次后进 DLQ, 循环有出口
坑 11 · 高峰期所有服务卡死, DB 锁等待 5000 — 症状: 全链雪崩在 prepare. 原因: 2PC prepare 后行锁被抱着等协调者发令, 协调者一慢全体锁表. 正解: 高并发换最终一致 (消息/补偿), 2PC 只留给低并发强一致场景。
# 错: 2PC prepare 后 5000 连接持锁等 commit 指令 → 全链雪崩
# 对: 拆本地事务 + Outbox, 锁持有时间从秒级降到毫秒级
坑 12 · 出了差异对不了账, 消息里只有一串 seq — 症状: 对账系统 join 不上. 原因: 消息体缺 orderNo/userId 等业务键, 无法与库内数据关联. 正解: 消息必带业务 ID + 事件时间 + 版本号。
# 错: body 只有 {"seq":881} → 对账系统 join 不上库表
# 对: body 带 orderNo/userId/amount/version/ts, 对账一键比对
坑 13 · 消费失败无限重试, 三分钟把商品服务打挂 — 症状: 消费侧变成机枪. 原因: 默认 requeue 或 retry 无限次, 毒消息持续射击下游. 正解: 有限次梯度重试 + 死信出口。
# 错: default-requeue-rejected: true(默认) + 手动 requeue 循环
# 对: retry.max-attempts: 3 + 进入 DLQ, 出口收敛
坑 14 · broker 主从切换, 发出去的 2000 条消息凭空消失 — 症状: 发送成功却查无此消息. 原因: RabbitMQ 未开 publisher confirm, 消息还在内存未落盘, 切换即丢. 正解: 开 confirm + 失败重投; Kafka 用 acks=all。
# 错: 未配置 confirm → broker 切主, 未落盘消息全丢
# 对: publisher-confirm-type: correlated + nack 时重投
坑 15 · 消息"已消费", 业务却没做 — 症状: offset 走了业务没走. 原因: auto.commit 每 5s 自动提交, 提交后处理才崩, 重启从下一条开始 — 丢处理不丢 offset. 正解: 关 auto.commit, 处理成功后再手动提交。
# 错: enable.auto.commit=true   # 5s 一提交, 崩在处理后 → 丢消息
# 对: enable.auto.commit=false + 处理成功后 commitSync()
坑 16 · 积压 400 万, 运维重启消费者集群, 积压翻倍 — 症状: 越救越火. 原因: 重启触发 group rebalance, 消费暂停 + 分区反复迁移, 还会收到 CommitFailedException. 正解: 先查 lag 定位原因, 扩容/跳过, 别乱重启。
# 错: systemctl restart consumer × 12 → rebalance 风暴, lag 再翻倍
# 对: kafka-consumer-groups.sh --describe 查 lag 定位, 再扩容/跳过
坑 17 · 加到 20 个消费者实例, 吞吐纹丝不动 — 症状: 扩容无效还白花钱. 原因: 12 个分区最多 12 个消费者同时消费, 多出的 8 个 assignment 为空, 空转. 正解: 实例数 = 分区数; 要更多并行先扩分区 (注意同 key 顺序前提)。
# 错: 12 分区 × 20 实例 → 8 个实例 assignment 为空, 空转
# 对: 实例数 = 分区数; 要更多并行 → 先 kafka-topics --alter 扩分区
坑 18 · 头部商户占了单分区 40% 流量, 消费者单核打满 — 症状: 顺序消息的数据倾斜. 原因: 按商户 key 分区保序, 头部效应把热 key 挤进一个分区. 正解: key 加订单级分片后缀, 全局有序降级片内有序。
// 错: key = merchantId            // → 头部商户独占分区
// 对: key = merchantId + "#" + (orderNo.hashCode() & 7)
坑 19 · Cancel 盲目解冻, 账户余额凭空变多 — 症状: TCC 资损. 原因: Try 没做资源预留检查、不记冻结流水, Cancel 不查流水盲目加钱, 重复 Cancel 双倍退款. 正解: Try 校验 + 冻结 + 记 xid 流水, Cancel 按 xid 幂等退冻结, 无流水即空回滚返回成功。
-- 错: cancel 直接 UPDATE bal = bal + ? WHERE uid=?  -- 双倍退款
-- 对: 按 xid 查冻结流水: 无→空回滚返回成功; 有→退冻结+流水置已退
坑 20 · 机房迁移后消费者集体"卡住不消费" — 症状: 新集群消费位置诡异. 原因: 新集群的分区 offset 与旧集群完全无关, 直接导入旧 offset 语义全错. 正解: MirrorMaker 2 迁数据, offset 翻译用其 checkpoint, 或按时间戳重置。
# 错: 把旧集群 __consumer_offsets 导入新集群 → offset 语义全错
# 对: kafka-consumer-groups.sh --reset-offsets --to-datetime <ts> --execute