学习目标
学完本文,你应该能够:
- 用一句话讲清分布式事务问题:一次大操作由多个落在不同服务/数据库的小操作组成,必须保证"要么同时成功,要么同时失败"。
- 说清 2PC 的两阶段 + 中段回滚流程,并指出它在互联网场景下落不了地的三大根因(死锁、性能、数据不一致)。
- 把"面试套路"练成肌肉记忆:先铺方案全景 → 再给可落地方案(MQ)→ 把 2PC/TCC 当交流引子,而不是一上来就背 2PC。
- 用 Go 跑通基于 MQ 的可靠消息投递:本地消息表 + 手动 Ack + 定时补偿,落地"最终一致性"。
- 讲清双向消息确认机制为什么是这类方案可落地的核心,并区分清楚"强一致"与"最终一致"的取舍边界。
前置知识:了解单库本地事务的 ACID;知道消息队列(RabbitMQ / RocketMQ / Kafka 任一即可)的基本收发模型;写过一点 Go(goroutine、channel、sync.Mutex 基础即可)。
本章你会动手做的事:
- 跑通一个 Go 写的 2PC 协调者/参与者模拟,亲眼看到"有一个参与者投 NO 就全局回滚"。
- 用 Go + 本地消息表实现一个"订单系统 ↔ 优惠券系统"的可靠消息投递,并看到补偿任务把丢失的消息捞回来。
- 画一张"分布式事务一致性"考点脑图(见下方),把 2PC / MQ / 双向确认串成一张网。
一、考点脑图
先给整篇文章一张总览图。分布式事务一致性这道题,考的是"你怎么在一致性、可用性、复杂度之间做工程取舍",而不是背某个协议。
二、先建立直觉:什么是分布式事务一致性
2.1 生活类比:部门团建 AA 付款
在互联网之前,一个交易系统把所有下单逻辑写在一个单体应用里,靠单库本地事务(一个 @Transactional)就能保证"全成或全败"。但系统拆分后,订单、商品、促销变成了三个独立服务、三套数据库——一个本地事务管不了三个库,分布式事务问题就出现了。
2.2 严谨定义
分布式事务:一次大的业务操作由多个小操作组成,这些小操作分别部署/存储在不同的服务器或数据库上;分布式事务要保证这些小操作要么同时成功,要么同时失败。
2.3 面试套路:别一上来就背 2PC
很多候选者被问到"怎么保证系统间的分布式一致性",会下意识选一种方案开讲:“可以基于两阶段提交……“然后啪啪讲 2PC 原理。这其实犯了一个明显错误——实际工作中我们很少用 2PC / 3PC / TCC,基本都是基于 MQ 的可靠消息投递。
然后给出可落地方案:互联网高并发场景基本都选"基于 MQ 的最终一致性",因为它解耦、削峰、可扩展。
最后把 2PC / TCC 当作交流引子——说明它们原理我懂、但工业界落地代价大,只适合金融支付这类强一致的一流场景。这样既能展开讨论,又显得我真的做过。
三、两阶段提交 2PC:教父级协议
3.1 原理:协调者与参与者
2PC 是分布式事务的"教父级"协议,最早来自数据库领域的 XA 规范(X/Open 提出)。它定义了**事务管理器(协调者)和资源管理器(参与者,如 MySQL / Oracle)**之间的接口。整个过程分两个阶段:
- 阶段一·准备(Prepare):协调者通知所有参与者"开启事务、你做好提交准备了吗?";参与者写
undo/redo日志、加资源锁,并返回 YES / NO。 - 阶段二·提交(Commit):只有当所有参与者都返回 YES,协调者才发
COMMIT;任意参与者返回 NO,协调者就发ROLLBACK,参与者用准备阶段记录的undo日志回滚。
sequenceDiagram
participant C as 协调者(事务管理器)
participant P1 as 订单库
participant P2 as 商品库
participant P3 as 促销库
Note over C: 阶段一 · 准备
C->>P1: 开启事务,能提交吗?
C->>P2: 开启事务,能提交吗?
C->>P3: 开启事务,能提交吗?
P1-->>C: YES(写 undo/redo 日志 + 加锁)
P2-->>C: YES
P3-->>C: NO(资源锁定失败)
Note over C: 阶段二 · 提交 / 中段回滚
C->>P1: ROLLBACK(用 undo 日志回滚)
C->>P2: ROLLBACK
C->>P3: ROLLBACK3.2 用 Go 跑通一个 2PC 模拟
下面这段 Go 代码把协调者和参与者都建模出来,真实可运行:准备阶段收集所有参与者的投票,只要有一个投 NO,就全局回滚。
package main
import "fmt"
// Participant 模拟一个资源管理器(如一个数据库分库)
type Participant struct {
name string
canCommit bool // 模拟该参与者是否能做好提交准备(资源能否锁定)
}
// prepare 阶段一:协调者询问参与者是否可以提交
// 真实场景里这里会写 undo/redo 日志并对数据行加锁
func (p *Participant) prepare() bool {
fmt.Printf("[准备] %s:记录 undo/redo 日志,返回 %v\n", p.name, p.canCommit)
return p.canCommit
}
func (p *Participant) commit() {
fmt.Printf("[提交] %s:正式写入数据并释放锁\n", p.name)
}
func (p *Participant) rollback() {
fmt.Printf("[回滚] %s:利用 undo 日志回滚并释放锁\n", p.name)
}
func main() {
participants := []*Participant{
{name: "订单库(数据库一)", canCommit: true},
{name: "商品库(数据库二)", canCommit: true},
{name: "促销库(数据库三)", canCommit: false}, // 模拟库存资源锁定失败
}
// 阶段一:准备,收集所有投票后再决策
allReady := true
for _, p := range participants {
if !p.prepare() {
allReady = false
}
}
// 阶段二:全 YES 则提交,否则中段回滚(投 YES 的也要跟着回滚)
if allReady {
for _, p := range participants {
p.commit()
}
fmt.Println("✅ 全局事务提交成功")
} else {
for _, p := range participants {
p.rollback()
}
fmt.Println("❌ 全局事务回滚(至少一个参与者未就绪)")
}
}
跑一下你会看到:促销库投了 NO,于是订单库、商品库即便准备成功,也被协调者要求回滚——这就是"要么全成、要么全败”。
3.3 2PC 为什么在互联网落不了地
2PC 能借助数据库本地事务"几乎不侵入业务"地实现一致性,但它的准备阶段必须加资源锁(如 MySQL 的行锁),由此带来三个致命问题:
② 性能低下:被锁的数据行,其他事务只能阻塞等待,分布式事务呈现高延迟、吞吐量低,根本扛不住海量并发。
③ 数据不一致:提交阶段协调者发 COMMIT 后若发生网络异常,只有部分库收到并执行,没收到的库永远不提交,系统出现不一致。
举个库存例子:库存=1,准备阶段问"能扣吗"回答"能”,但不锁行的话,提交前另一个请求把库存扣成 0,等你提交阶段再去扣,库存就变成 -1 了。所以必须锁——但一锁,上面三个问题就全来了。
也正因如此,互联网几乎不用 2PC,而是改用下面要讲的 MQ 方案。
四、为什么互联网选 MQ 可靠消息投递
4.1 思路:放弃强一致,拥抱最终一致
应对高并发,工业界的主流做法是放弃强一致性、选择最终一致性,用消息队列把"同步阻塞的三方协调"变成"异步解耦的点对点"。还是以下单为例:
订单系统不直接同步调用优惠券系统,而是把"扣减优惠券"这件事,作为一条已持久化的消息放进 MQ,由优惠券系统异步消费执行。只要这条消息最终能在优惠券系统里被执行,一致性就达成了。
sequenceDiagram
participant O as 订单系统
participant MQ as 消息队列
participant C as 优惠券系统
O->>MQ: 投递“扣减优惠券”消息(持久化)
Note over MQ: 消息落盘,宕机重启也不丢
MQ->>C: 推送消息
C->>C: 扣减优惠券(本地事务)这样做一举三得:
- 解同步阻塞:订单系统投完消息即可返回,不用等优惠券系统处理完。
- 业务解耦:订单系统和优惠券系统互不依赖,各自独立演进、独立扩容。
- 流量削峰:大促瞬时流量先堆在 MQ 里,优惠券系统按自己节奏消费。
4.2 坑一:MQ 自动应答导致消息丢失
这是面试官最爱追问的点。以优惠券系统消费为例:MQ 默认开启自动应答(autoAck)——消费者一收到消息,MQ 就立刻把这条持久化消息删了。可优惠券系统执行过程中一旦抛异常中断,消息就没了,扣券永远没发生,消息丢失。
下面是一段贴近 RabbitMQ 的手工 Ack 写法(核心在于 autoAck=false + 成功后才 Ack):
// 订阅时务必关闭自动应答
msgs, _ := ch.Consume(queue, consumer, /* autoAck = */ false, false, false, false, nil)
for d := range msgs {
// 1) 先执行业务:扣减优惠券
err := deductCoupon(context.Background(), d.Body)
if err != nil {
// 2) 业务失败:不 Ack,MQ 会在重试策略下重新投递
// 超过最大重试次数会进死信队列,等待人工干预
_ = d.Nack(false, true) // multiple=false, requeue=true
continue
}
// 3) 业务成功“之后”才手动 Ack,MQ 才真正删除消息
if err := d.Ack(false); err != nil {
// Ack 本身失败也要记录告警,避免消息静默丢失
log.Printf("ack failed: %v", err)
}
}
4.3 坑二:消息积压与死信队列
大促瞬时流量剧增,大量消息来不及消费、积压在 MQ。若优惠券系统因限流等原因长时间消费不动,消息会被 MQ 不断重试,超过最大重试次数后丢弃进死信队列(DLQ)——而进死信的消息往往需要人工干预,实际大概率被"静默丢弃",造成一致性缺口。
五、可落地的核心:双向消息确认机制
5.1 为什么需要"双向确认"
订单系统投出消息后,作为生产者它并不知道优惠券系统(消费者)是成功还是失败。如果让订单系统能感知消费响应,即使 MQ 把消息弄丢了,订单系统也能通过定时任务扫描,把未完成的消息重新投递——这就是双向消息确认,也是基于 MQ 实现分布式事务可落地的关键。
5.2 落地流程
- 订单系统把要发的消息先持久化到本地消息表,状态置为「待发送(PENDING)」(与下单业务在同一本地事务内落库)。
- 订单系统把消息投递到 MQ。
- 优惠券系统消费成功,向 MQ 回发一条确认消息。
- 订单系统收到确认,把本地消息表里的该条记录状态改为「已完成(DONE)」。
- 定时任务扫描一段时间内仍处于「待发送/已发送」状态的消息,重新投递,完成补偿。
sequenceDiagram
participant O as 订单系统
participant DB as 本地消息表
participant MQ as 消息队列
participant C as 优惠券系统
O->>DB: 1. 下单事务内写消息(状态=待发送)
O->>MQ: 2. 投递消息
MQ->>C: 3. 推送扣券消息
C->>C: 4. 扣减优惠券(业务执行)
C->>MQ: 5. 消费成功,回发确认
MQ->>O: 6. 确认通知
O->>DB: 7. 更新消息状态=已完成
Note over O,DB: 补偿:定时任务扫描未完成消息,重新投递5.3 用 Go 跑通本地消息表 + 补偿
下面这段 Go 程序把上面的流程全部落到了代码,真实可运行:OrderDB 就是本地消息表,正常流程里 O1001 / O1002 顺利消费;O1003 的消息在投递时"丢失"(没进 MQ),最后补偿任务把它捞出来重新投递。
package main
import (
"fmt"
"sync"
"time"
)
// ---- 本地消息表 ----
type MsgStatus string
const (
StatusPending MsgStatus = "PENDING" // 待发送
StatusSent MsgStatus = "SENT" // 已投递
StatusDone MsgStatus = "DONE" // 已完成(消费者已确认)
)
type LocalMessage struct {
ID string
BizKey string // 业务键,如 order_id
Payload string
Status MsgStatus
}
// OrderDB 扮演订单系统侧的“本地消息表”
type OrderDB struct {
mu sync.Mutex
messages map[string]*LocalMessage
}
func NewOrderDB() *OrderDB { return &OrderDB{messages: make(map[string]*LocalMessage)} }
// 步骤1:下单时,业务与消息在同一个本地事务里落库
func (db *OrderDB) createOrderWithMsg(orderID, payload string) {
db.mu.Lock()
defer db.mu.Unlock()
db.messages[orderID] = &LocalMessage{
ID: orderID, BizKey: orderID, Payload: payload, Status: StatusPending,
}
fmt.Printf("[订单系统] 本地事务落库:order=%s, 状态=%s\n", orderID, StatusPending)
}
// 步骤2:投递消息到 MQ(此处用打印模拟)
func (db *OrderDB) deliver(orderID string) {
db.mu.Lock()
m := db.messages[orderID]
if m != nil && (m.Status == StatusPending || m.Status == StatusSent) {
m.Status = StatusSent
}
db.mu.Unlock()
if m != nil {
fmt.Printf("[订单系统] 投递 MQ:order=%s (状态:%s)\n", orderID, m.Status)
}
}
// 步骤4:收到消费者确认,标记完成
func (db *OrderDB) confirm(orderID string) {
db.mu.Lock()
defer db.mu.Unlock()
if m, ok := db.messages[orderID]; ok {
m.Status = StatusDone
fmt.Printf("[订单系统] 收到消费确认,标记完成:order=%s\n", orderID)
}
}
// 步骤5:定时任务扫描未完成消息,重新投递(补偿)
func (db *OrderDB) resendPending() {
db.mu.Lock()
var pending []string
for id, m := range db.messages {
if m.Status == StatusPending || m.Status == StatusSent {
pending = append(pending, id)
}
}
db.mu.Unlock()
for _, id := range pending {
fmt.Printf("[补偿任务] 发现未完成消息,重新投递:order=%s\n", id)
db.deliver(id)
}
}
// ---- 优惠券系统侧 ----
// consume 消费消息,成功后回发确认(手动 Ack 思想)
func consume(mq chan string, db *OrderDB, wg *sync.WaitGroup) {
defer wg.Done()
for orderID := range mq {
// 模拟扣减优惠券成功
fmt.Printf("[优惠券系统] 扣减优惠券成功:order=%s\n", orderID)
// 关键:业务成功之后才确认(对应 MQ 的手动 Ack)
db.confirm(orderID)
}
}
func main() {
db := NewOrderDB()
mq := make(chan string, 10)
var wg sync.WaitGroup
wg.Add(1)
go consume(mq, db, &wg)
// 正常下单:O1001 / O1002 消息顺利送达优惠券系统
for _, o := range []string{"O1001", "O1002"} {
db.createOrderWithMsg(o, "deduct-coupon")
db.deliver(o)
mq <- o
}
// 模拟 MQ 网络异常:O1003 的扣券消息在投递时丢失(未进 MQ)
db.createOrderWithMsg("O1003", "deduct-coupon")
// 注意:这里没有 db.deliver("O1003"),也没有 mq <- "O1003"
close(mq)
wg.Wait()
// 补偿任务:扫描本地消息表中未完成的消息,重新投递
db.resendPending()
time.Sleep(100 * time.Millisecond)
}
运行后你会看到:O1001 / O1002 一路走到"标记完成",而 O1003 因为消息丢失一直停在 PENDING,最后被补偿任务捞出来重新投递——只要本地消息表还在,消息就丢不了。
5.4 生产化的本地消息表
上面用 map 演示了逻辑。真实项目里本地消息表是一张物理表,与业务订单同一事务落库,确保"订单成了、消息也一定在":
CREATE TABLE local_message (
id BIGSERIAL PRIMARY KEY,
biz_key VARCHAR(64) NOT NULL,
payload TEXT NOT NULL,
status VARCHAR(16) NOT NULL DEFAULT 'PENDING', -- PENDING / SENT / DONE
created_at TIMESTAMP NOT NULL DEFAULT NOW(),
updated_at TIMESTAMP NOT NULL DEFAULT NOW()
);
CREATE INDEX idx_status ON local_message (status, created_at); -- 补偿扫描靠它
六、延伸:TCC 与 2PC 有什么不同(思考题)
课程留了一个思考题:还有一种叫 TCC(Try-Confirm-Cancel)的方案,它和 2PC 不同点在哪里?这里先给个对比框架,方便你面试时展开:
- 锁的粒度与时长:2PC 在准备阶段就加数据库行锁直到提交/回滚,锁持有时间长、易阻塞;TCC 把事务拆成
Try(预留资源,如冻结库存而不是真实扣减)、Confirm(真正提交)、Cancel(释放预留),不长期持有数据库锁,靠业务层的预留/补偿来实现。 - 侵入性:2PC 由数据库/XA 层托管,业务侵入小;TCC 要手写三个接口(Try/Confirm/Cancel),业务侵入大、开发成本高。
- 适用面:2PC 偏数据库层面、适合强一致但低并发;TCC 偏业务层面、能扛更高并发,但只适合少数强一致场景(如金融转账),落地代价依然不小——所以互联网主流仍是 MQ 最终一致。
七、自测题与动手练习
下面几道题专门用来检验你是"真懂"还是"只会背"。
单库本地事务(一个
@Transactional)只能管住一个数据库实例内的一组操作;一旦订单、商品、促销拆成三个独立库,一个本地事务就管不到另外两库了,所以需要跨库的协调方案。不用它是因为准备阶段必须加数据库行锁,会带来三大问题:死锁(故障后资源锁死、数据库阻塞)、性能低下(锁住的住行其他事务只能等)、数据不一致(提交阶段网络异常,部分库收到 COMMIT、部分没收到)。所以扛不住海量并发,互联网基本不用。
核心是双向消息确认——订单系统(生产者)也能感知消费者有没有成功消费,配合定时任务扫描未完成的消息重新投递做补偿。这套方案解耦、削峰,天然适合高并发。
防法是关掉自动应答、改手动 Ack:只有业务逻辑真正执行成功之后,才向 MQ 发 Ack,MQ 才删消息;失败就用 Nack 让 MQ 重投,超次数进死信队列等人工处理。这种细节最能体现“真做过”。
再加上对账任务(离线比对订单与优惠券系统的状态)做最后一道防线,彻底兜底。
动手练习(建议真做一遍):
- 把本文第 3.2 节的 2PC 模拟跑起来,把第三个参与者的
canCommit改成true,看输出如何变成"全局提交成功"。 - 把第 5.3 节的本地消息表程序跑通,再故意把
consume里的confirm注释掉,观察补偿任务扫描到多少条"未完成"消息。 - 用一张 SQL 本地消息表 + 你熟悉的 MQ(RabbitMQ / RocketMQ 任一),把"订单 ↔ 优惠券"的双向确认真实落一遍,并写一段选型分析:为什么这条链路选最终一致而不是 2PC。
八、本章小结
- 分布式事务 = 跨服务/跨库的"全成或全败":单库本地事务管不到多库,所以才需要协调方案。
- 2PC 是教父级协议但落不了地:准备阶段加行锁,导致死锁、性能低、数据不一致三大问题——面试用它当"交流引子",别当"答案"。
- 工业界真正用的是 MQ 可靠消息投递(最终一致):解耦 + 削峰 + 可扩展;防消息丢失靠手动 Ack,防积压丢消息靠双向确认 + 定时补偿。
- 可落地性的核心只有一句:实现生产者↔消费者的双向确认;实际工作中并非所有业务都要强一致,站在业务场景权衡成本才是高手做派。
- 下一篇可以深入 TCC / Saga 这类补偿型事务,看看在"不能丢、但要高并发"的金融场景里,工程上怎么把一致性"算"出来。