跨服务的一致性拼图: 能不用分布式事务就不用 — 本地事务保原子, 消息保可达, 幂等保不重, 补偿保回退, 对账保最终兜底
把跨服务一致性想成两个部门的报销流程: 2PC 是"两边同时签字才生效", 任何一边卡住全体干等; TCC 是"先占额度, 再正式划扣, 不行就释放"; Saga 是"先各办各的, 出了岔子按流程冲销"; Outbox 则是"把要传的话写进同一个小本本, 专人负责抄送"。
机制本质: 分布式事务的终极答案不是"更强的 2PC", 而是把长事务拆短 — 本地事务保原子, 消息保可达, 幂等保不重, 补偿保回退, 对账保最终兜底。能最终一致就别强一致: 一致性是有价格的, 且价格随并发指数上涨。
BEGIN; UPDATE account SET bal = bal - 50 WHERE uid = 1; UPDATE account SET bal = bal + 50 WHERE uid = 2; -- 关键: 同库才能同事务 COMMIT; -- → 跨了库/服务, 这两条就没人替你保证同生共死
阶段1: coordinator → 全体: prepare? 参与者锁行、写 redo, 回 yes 阶段2: 全票 yes → commit; 任一 no → rollback -- 关键: 协调者崩在两阶段之间, 参与者抱着锁干等 -- → 协调者单点 + 同步阻塞, 高并发场景基本弃用
2PC: prepare ── commit (参与者无超时, 卡死等协调者) 3PC: canCommit ── preCommit ── doCommit (超时可自行推进) -- 关键: 用超时换掉阻塞, 但分区下仍可能双提交
@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 退冻结
正向: T1 创建订单 → T2 扣款 → T3 扣库存 T3 失败 ⇒ C2 退款 → C1 关闭订单 -- 关键: 补偿必须可重试且幂等, 失败转人工干预 -- → 最终回滚到 T1 之前的业务等价状态
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
INSERT INTO inbox(msg_id, consumed_at) VALUES('evt_1', NOW()); -- → Duplicate entry 'evt_1' for key 'PRIMARY' -- 关键: 唯一键冲突 = 已消费, 捕获后直接 ACK 跳过
生产者 → broker: 未收到 ack → 重发 → 消费者收到 2 次
消费者 → broker: 处理完未 ack 就崩 → rebalance → 再投 1 次
# 关键: at-least-once 是"不丢"的代价, 重投无法根除只能去重enable.idempotence=true # PID + seq: 单分区内不重不乱 transactional.id=svc-1 # 跨分区原子 + 崩溃后 fencing 旧实例 # 关键: consumer 端 isolation.level=read_committed 才配套生效
ALTER TABLE payment ADD UNIQUE KEY uk_order_no (order_no); -- 重复请求 INSERT → 1062 Duplicate entry 'PO20260901' -- 关键: 幂等键用业务自然键, 自增 ID 每次都不同挡不住
-- 先 Redis: SETNX msg:m_88 EX 86400 挡洪峰, DB 唯一键终审 INSERT INTO dedup(msg_id, expire_at) VALUES('m_88', NOW() + INTERVAL 1 DAY); -- 关键: 去重记录必须带过期清理, 否则 90 天后慢查询拖垮库
失败 → retry.5s → retry.30s → retry.5m → DLQ # 关键: 每级都设 TTL + 死信路由, 到点自动晋级, 不死循环 # → 第 4 级还失败就进 DLQ, 等人工/定时修复重放
QueueBuilder.durable("order.work") .withArgument("x-dead-letter-exchange", "order.dlx") // 死信去向 .withArgument("x-message-ttl", 60000) // 60s 后进 DLX .build(); // 关键: DLQ 深度告警 + 修复重放脚本先备好
# 同一订单的消息带同 key → hash 落同一分区, 分区内 FIFO producer.send(new ProducerRecord("orders", orderId.getBytes(), payload)); enable.idempotence=true # 关键: in-flight≤5 时失败重试也不乱序
之前"先写订单再发 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=? -- 关键: 投递失败只标记不删除, 下轮继续 — 事件永不丢
要求"订单失败库存一定不扣"。先发半消息, 执行本地事务, 成功 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); // 关键: 半消息先行, 回查兜底
消费 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, 才看不到未提交结果
所有消费者按同一模板写: 去重表 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; -- 关键: 去重记录和业务结果同事务, 崩了一起回滚, 重投再来
毒消息曾把消费者 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
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 退冻结, 幂等可重试 }
没设 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) 比较版本, 只前滚不回滚
上游 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 留证 # 关键: 先止血(跳过)再补数(离线回放), 别让正常业务陪葬
不动业务代码: 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 把"已提交的行变更"变成事件流, 天然无双写分叉
按商户 ID 分区保序, 结果头部商户独占 40% 流量把单分区打满。打散: 商户内按订单再分片, 全局有序降级为片内有序。
// 错: key = merchantId // 头部商户占该分区 40% 流量, 消费者单核打满 // 对: 分片 = orderNo.hashCode() & 7 // 商户内按订单拆 8 个逻辑片 // key = merchantId + "#" + 分片 // 全局有序降级为片内有序: 同一订单仍严格有序, 不同订单可并行消费 // 监控: 每分区消息数方差 < 20% 才算打散成功 // 关键: 片数取 2 的幂, 位运算代替取模, 分片更均匀
-- 错: COMMIT 订单 → send(mq) 失败 → 库存永远不知道这一单 -- 对: 同事务 INSERT outbox → relay 异步投递, 崩了也不丢事件 -- 消费端再配 inbox 去重, 重投也不重做
# 错: enable.idempotence=true 就以为万无一失 → 消费写 DB 仍重复 # 对: DB 侧 inbox 去重表 + 唯一键, 消费端幂等自己做 # EOS 只管流内部, 出口是另一个系统的事
# 错: T1 建单 → T2 扣款 → T3 扣库存(失败) → 无 C2 可走, 卡死 # 对: T 配 C 成对设计: C1 关单 / C2 退款 / C3 回补库存 # 上线前评审表: 每个 T 的 C 是什么, 写进设计文档
-- 错: if !refund() { log.error("cancel fail") } -- 沉没在日志里 -- 对: INSERT INTO manual_task(xid, action, retry_cnt) -- -- 关键: 补偿失败必须升级为人 + 定时重推
# 错: 配了 DLX 就不管 → 30w 条死了 3 个月才发现 # 对: 深度>100 告警; 修复后按原 exchange 重新 publish 回业务队列 # 每周固定重放窗口, 重放前先修数据再回放
// 错: new ProducerRecord("orders", payload) // 无 key 轮询 // 对: new ProducerRecord("orders", orderId, payload) // 关键: 同 key 同分区
# 错: message.max.bytes=10485760 # 调到 10MB, 复制带宽翻 10 倍 # 对: 消息体只放 {"fileUrl":"oss://..."}, 文件走对象存储 # 消费端按 URL 拉取, 失败可重试且不占 broker
// 错: 永远 return UNKNOW → 半消息悬挂, 3 次回查后丢弃 // 对: 查库有单→COMMIT; 无单→回查计数<3 return UNKNOW, 之后 ROLLBACK
-- 错: id BIGINT AUTO_INCREMENT 做唯一键 → 重投 100 次插 100 行 -- 对: UNIQUE KEY uk_biz (order_no) -- 关键: 业务键才拦得住重复
# 错: 失败 requeue 回原队列 → 毒消息死循环, CPU 100% # 对: retry 队列 TTL 5m → 3 次后进 DLQ, 循环有出口
# 错: 2PC prepare 后 5000 连接持锁等 commit 指令 → 全链雪崩 # 对: 拆本地事务 + Outbox, 锁持有时间从秒级降到毫秒级
# 错: body 只有 {"seq":881} → 对账系统 join 不上库表 # 对: body 带 orderNo/userId/amount/version/ts, 对账一键比对
# 错: default-requeue-rejected: true(默认) + 手动 requeue 循环 # 对: retry.max-attempts: 3 + 进入 DLQ, 出口收敛
# 错: 未配置 confirm → broker 切主, 未落盘消息全丢 # 对: publisher-confirm-type: correlated + nack 时重投
# 错: enable.auto.commit=true # 5s 一提交, 崩在处理后 → 丢消息 # 对: enable.auto.commit=false + 处理成功后 commitSync()
# 错: systemctl restart consumer × 12 → rebalance 风暴, lag 再翻倍 # 对: kafka-consumer-groups.sh --describe 查 lag 定位, 再扩容/跳过
# 错: 12 分区 × 20 实例 → 8 个实例 assignment 为空, 空转 # 对: 实例数 = 分区数; 要更多并行 → 先 kafka-topics --alter 扩分区
// 错: key = merchantId // → 头部商户独占分区 // 对: key = merchantId + "#" + (orderNo.hashCode() & 7)
-- 错: cancel 直接 UPDATE bal = bal + ? WHERE uid=? -- 双倍退款 -- 对: 按 xid 查冻结流水: 无→空回滚返回成功; 有→退冻结+流水置已退
# 错: 把旧集群 __consumer_offsets 导入新集群 → offset 语义全错 # 对: kafka-consumer-groups.sh --reset-offsets --to-datetime <ts> --execute