分布式事务怎么搞?Go后端新人也能看懂的最佳实践
大家好,我是你们的学长,一名211高校计算机专业的研二学生。平时喜欢写技术博客,尤其热衷于把复杂的技术概念“翻译”成大白话,帮助刚入门的同学们少走弯路。今天我想和大家聊聊一个听起来很高大上、但其实在很多实际项目中都会遇到的问题:分布式事务。
我当初第一次听到“分布式事务”这个词时,脑子里全是问号:事务不就是数据库里那个 ACID 吗?怎么还分“分布”和“不分布”?后来在实习做电商订单系统时,才真正体会到它的痛——用户付了钱,但库存没扣,或者优惠券用了却没生效……这些看似小问题,放到线上就是严重事故。
所以今天这篇教程,我会用最直白的语言,配合 Go 语言的实战代码,带你从零开始理解分布式事务的核心思想,并掌握几种主流的解决方案。无论你是刚学完 Go 基础的新手,还是正在准备后端岗位面试的同学(对,这可是高频面试题挑战!),都能从中受益。
一、什么是分布式事务?为什么需要它?
单体 vs 分布式:问题从哪来?
想象一下你写了一个简单的电商下单功能:
- 用户点击“购买”
- 系统扣减库存
- 创建订单
- 扣除用户余额
如果所有操作都在同一个数据库里完成,那很简单:用一个 BEGIN...COMMIT 包起来就行。这就是传统意义上的本地事务,满足 ACIC 特性(原子性、一致性、隔离性、持久性)。
但现在,你的系统变大了:
- 库存服务 → 用 MySQL A
- 订单服务 → 用 MySQL B
- 账户服务 → 用 PostgreSQL C
三个服务,三个数据库。这时,如果“扣库存”成功了,但“创建订单”失败了,怎么办?总不能让用户白拿商品吧!
分布式事务要解决的就是:跨多个数据库或服务的操作,如何保证要么全部成功,要么全部失败?
💡 小贴士:只要你的系统涉及多个数据源(哪怕是同一个数据库的不同表,但由不同微服务管理),就可能遇到分布式事务问题。
二、开发环境准备(Go + Docker 快速搭建)
为了让你能动手实践,我们先搭个最小可用环境。
1. 安装必备工具
- Go 1.20+
- Docker + Docker Compose(用于快速启动多个数据库)
- 任意代码编辑器(推荐 VS Code)
2. 创建项目结构
mkdir distributed-tx-demo
cd distributed-tx-demo
go mod init distributed-tx-demo
3. 用 Docker 启动两个数据库
创建 docker-compose.yml:
version: '3'
services:
mysql-order:
image: mysql:8.0
container_name: mysql-order
environment:
MYSQL_ROOT_PASSWORD: root
MYSQL_DATABASE: order_db
ports:
- "3307:3306"
mysql-inventory:
image: mysql:8.0
container_name: mysql-inventory
environment:
MYSQL_ROOT_PASSWORD: root
MYSQL_DATABASE: inventory_db
ports:
- "3308:3306"
运行:
docker-compose up -d
现在你有两个独立的 MySQL 实例:
- 订单库:localhost:3307
- 库存库:localhost:3308
4. 安装 Go 依赖
go get github.com/go-sql-driver/mysql
go get github.com/jmoiron/sqlx
三、核心方案对比:CAP、BASE 与三大主流解法
在深入代码前,先了解几个关键理论。
CAP 定理(别怕,很简单)
一个分布式系统最多只能同时满足以下两点:
- C(Consistency):数据一致
- A(Availability):服务可用
- P(Partition tolerance):网络分区容忍
现实中,P 是必须的(网络可能断),所以只能在 C 和 A 之间权衡。分布式事务本质是在牺牲部分强一致性,换取可用性和性能。
BASE 理论(最终一致性)
- Basically Available(基本可用)
- Soft state(软状态)
- Eventual consistency(最终一致性)
大多数分布式事务方案都基于 BASE,而不是强 ACID。
三大主流解决方案对比
| 方案 | 原理 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| 2PC(两阶段提交) | 协调者统一指挥各参与者准备→提交 | 强一致性 | 性能差、同步阻塞、单点故障 | 传统金融系统(如银行) |
| TCC(Try-Confirm-Cancel) | 业务层面实现补偿逻辑 | 灵活、性能好 | 开发复杂度高 | 高并发电商、支付 |
| 消息队列+本地事务表 | 用消息确保异步最终一致 | 简单、解耦 | 有延迟、需处理消息重复 | 订单通知、积分发放 |
📌 作为新手,建议优先掌握 消息队列方案 和 TCC 的思想,2PC 了解即可(面试常问,但生产少用)。
四、实战:用“本地消息表”实现最终一致性(Go 示例)
我们来做一个简化版的“下单扣库存”流程。
步骤设计(文字流程图)
1. 用户下单
2. 在订单库中:
- 插入订单记录(状态=待支付)
- 插入一条“扣库存消息”到本地消息表(状态=待发送)
3. 启动后台任务,扫描本地消息表
4. 发送消息到 MQ(比如 RabbitMQ 或 Kafka)
5. 库存服务消费消息,执行扣库存
6. 成功后,更新本地消息表状态为“已发送”
✅ 关键点:“插入订单”和“插入消息”在同一个本地事务中,保证原子性。
1. 创建数据库表
订单库(order_db)
CREATE TABLE orders (
id BIGINT PRIMARY KEY,
user_id BIGINT,
product_id BIGINT,
status VARCHAR(20)
);
CREATE TABLE outbox_messages (
id BIGINT AUTO_INCREMENT,
message_type VARCHAR(50),
payload JSON,
status ENUM('pending', 'sent') DEFAULT 'pending',
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (id)
);
库存库(inventory_db)
CREATE TABLE inventory (
product_id BIGINT PRIMARY KEY,
stock INT
);
-- 初始化商品库存
INSERT INTO inventory (product_id, stock) VALUES (1001, 100);
2. Go 代码实现
连接数据库
// db.go
package main
import (
"database/sql"
"log"
_ "github.com/go-sql-driver/mysql"
"github.com/jmoiron/sqlx"
)
var orderDB *sqlx.DB
var inventoryDB *sqlx.DB
func init() {
var err error
orderDB, err = sqlx.Open("mysql", "root:root@tcp(localhost:3307)/order_db")
if err != nil {
log.Fatal(err)
}
inventoryDB, err = sqlx.Open("mysql", "root:root@tcp(localhost:3308)/inventory_db")
if err != nil {
log.Fatal(err)
}
}
下单 + 写消息(关键!)
// order_service.go
package main
import (
"context"
"database/sql"
"encoding/json"
"log"
)
type Order struct {
ID int64 `db:"id"`
UserID int64 `db:"user_id"`
ProductID int64 `db:"product_id"`
Status string `db:"status"`
}
type InventoryMessage struct {
OrderID int64 `json:"order_id"`
ProductID int64 `json:"product_id"`
Quantity int `json:"quantity"`
}
func CreateOrder(ctx context.Context, userID, productID int64) error {
tx, err := orderDB.BeginTxx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
// 1. 创建订单
_, err = tx.NamedExec(`
INSERT INTO orders (id, user_id, product_id, status)
VALUES (:id, :user_id, :product_id, :status)`,
Order{
ID: 12345, // 简化,实际用雪花ID
UserID: userID,
ProductID: productID,
Status: "created",
})
if err != nil {
return err
}
// 2. 写本地消息表(同一事务!)
msg := InventoryMessage{
OrderID: 12345,
ProductID: productID,
Quantity: 1,
}
payload, _ := json.Marshal(msg)
_, err = tx.Exec(`
INSERT INTO outbox_messages (message_type, payload)
VALUES (?, ?)`,
"DECREASE_INVENTORY", payload)
if err != nil {
return err
}
return tx.Commit()
}
消息发送器(模拟)
// message_sender.go
package main
import (
"context"
"database/sql"
"encoding/json"
"log"
"time"
)
// 模拟发送到 MQ(这里直接调库存服务)
func StartMessageSender() {
ticker := time.NewTicker(5 * time.Second)
for range ticker.C {
sendPendingMessages()
}
}
func sendPendingMessages() {
rows, err := orderDB.Query(`
SELECT id, payload FROM outbox_messages
WHERE status = 'pending' LIMIT 10`)
if err != nil {
log.Println("query messages failed:", err)
return
}
defer rows.Close()
for rows.Next() {
var id int64
var payload []byte
if err := rows.Scan(&id, &payload); err != nil {
continue
}
var msg InventoryMessage
if err := json.Unmarshal(payload, &msg); err != nil {
continue
}
// 调用库存服务(简化:直接操作库存DB)
if err := decreaseInventory(msg.ProductID, msg.Quantity); err == nil {
// 标记消息为已发送
orderDB.Exec("UPDATE outbox_messages SET status='sent' WHERE id=?", id)
log.Printf("Sent message %d for order %d", id, msg.OrderID)
} else {
log.Printf("Failed to process message %d: %v", id, err)
// 可重试,或告警
}
}
}
func decreaseInventory(productID int64, quantity int) error {
result, err := inventoryDB.Exec(`
UPDATE inventory
SET stock = stock - ?
WHERE product_id = ? AND stock >= ?`,
quantity, productID, quantity)
if err != nil {
return err
}
affected, _ := result.RowsAffected()
if affected == 0 {
return sql.ErrNoRows // 库存不足
}
return nil
}
主函数
// main.go
package main
import (
"context"
"log"
"time"
)
func main() {
// 模拟用户下单
go func() {
time.Sleep(2 * time.Second)
if err := CreateOrder(context.Background(), 101, 1001); err != nil {
log.Fatal("Create order failed:", err)
}
log.Println("Order created!")
}()
// 启动消息发送器
StartMessageSender()
}
运行后,你会看到:
- 订单被创建
- 5秒内,库存被扣减
- 消息状态变为“sent”
✅ 即使库存服务暂时不可用,消息也会在下一轮被重试,最终达成一致!
五、新手常见问题解答
Q1:消息重复消费怎么办?
A:幂等性是关键!在库存服务中,你可以:
- 为每条消息生成唯一 ID
- 在库存库中加一张
processed_messages表,记录已处理的消息 ID - 执行前先查是否处理过
-- 库存服务伪代码
IF NOT EXISTS (SELECT 1 FROM processed_messages WHERE message_id = ?) THEN
扣库存;
INSERT INTO processed_messages VALUES (?);
END IF;
Q2:本地消息表会不会越来越大?
A:会!但你可以:
- 定期归档或删除
status='sent'且超过7天的消息 - 使用单独的消息表,与业务表解耦
Q3:为什么不用 RocketMQ 的事务消息?
A:那是进阶方案!本地消息表更简单、通用,适合学习。RocketMQ/Kafka 的事务消息本质也是类似思想,只是把“本地消息表”内置到 MQ 里了。
Q4:TCC 怎么写?能给个例子吗?
A:TCC 分三步:
- Try:冻结资源(如预扣库存)
- Confirm:真正扣减
- Cancel:释放冻结
Go 中你需要为每个服务实现这三个接口。比如库存服务:
type InventoryService interface {
TryDecrease(productID int64, qty int) error
ConfirmDecrease(orderID int64) error
CancelDecrease(orderID int64) error
}
但 TCC 对业务侵入大,新手容易出错,建议先掌握本地消息表。
六、学习建议与避坑指南
下一步学什么?
- 深入消息队列:学 RabbitMQ / Kafka 的 Exactly-Once 语义
- Seata 框架:阿里开源的分布式事务中间件,支持 AT/TCC/Saga 模式
- Saga 模式:长事务的另一种解法,适合业务流程长的场景
面试题挑战准备
- “分布式事务有哪些方案?优缺点?”
- “2PC 和 3PC 的区别?”
- “如何保证消息不丢失、不重复?”
- “最终一致性 vs 强一致性,如何选择?”
💬 我面试时就被问过:“如果本地消息表写成功,但服务崩溃了,消息还能发出去吗?”
答案是:能! 因为后台任务会定期扫描未发送的消息,这就是“可靠事件模式”的核心。
避坑提醒
- ❌ 不要一上来就用 2PC —— 性能瓶颈大
- ✅ 优先考虑“业务可接受最终一致性”的场景
- ✅ 所有远程调用都要考虑超时、重试、熔断
- ✅ 日志要打全!方便排查“到底哪一步失败了”
结语
分布式事务听起来吓人,但拆开来看,无非是用各种手段保证多步操作的可靠性。作为后端开发者,你不需要一开始就精通所有方案,但一定要理解为什么需要它,以及最常用的“本地消息表”是如何工作的。
希望这篇技术分享能帮你迈出第一步。如果你觉得有用,欢迎关注我的博客,后续我会继续用 Go 带你攻克更多后端难题!
最后送大家一句话:“分布式系统的坑,都是踩出来的;但别人的坑,你可以绕着走。”
加油,未来的后端工程师!

评论 0