Go 基础体系 · 第 55/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。
Go 分布式事务实践:本地事务、Outbox、Saga、TCC 与 DTM
本文以 Go 1.26.4 和经 Go 模块版本列表核对的稳定版 DTM Go client v1.18.7 为基准。分布式事务不是把 BEGIN 扩展到网络,而是在进程崩溃、消息重复、请求超时和部分服务不可用时,仍让多个本地状态最终收敛。最重要的设计结果不是“所有步骤看起来同时成功”,而是每个中间状态可识别、每次重试可判定、每个失败可恢复。
本文用订单支付说明数据与调用生命周期。示例采用关系数据库和 Go 标准库表达核心机制;数据库错误码、锁语义和消息代理投递保证必须在实际选型上验证。
1. 先判断是否真的需要分布式事务
若订单、库存和余额能放进同一个数据库实例,一个短小的本地事务通常最可靠。拆库只因为“以后可能扩展”会立刻引入重复、乱序、补偿和对账成本。只有服务需要独立部署、数据受不同团队或合规边界管理,或单库容量已经有证据不足时,才值得跨边界协调。
一致性目标也要具体化:付款后库存必须立即可见,还是五秒内收敛即可;失败时是自动退款,还是订单进入人工审核;用户能否重复提交;对账以哪个系统为准。CAP 口号不能替代业务不变量。常见选择如下:
- 单库事务:最强且最简单,优先使用;
- Transactional Outbox:本地写与事件原子,消费者最终一致;
- Saga:长流程由多个本地事务及补偿组成;
- TCC:先预留资源,再确认或释放,实时性更强但接口成本高;
- XA/两阶段提交:参与者支持时可强协调,但可用性、锁持有和运维成本高。
2. 网络会制造三种结果,而不只是成功和失败
调用者收到成功,可以确认远端完成;收到明确业务拒绝,可以确认没有执行;超时、断线或网关重启则产生“结果未知”。远端可能没收到、执行到一半、已经提交但响应丢失。把 timeout 直接解释成失败并再次扣款,会造成双扣。
order -> payment: POST /charges (idempotency-key=pay-42)
payment -> database: COMMIT
payment -x-> order: 200 OK 响应丢失
order: 只能记录 unknown,并按 pay-42 查询或安全重试
因此每个写命令需要稳定业务键,如 order_id + operation 或客户端生成的 UUID。服务端在同一事务中保存键和结果;重复请求返回原结果,而不是重新执行。超时预算应从入口 context 贯穿数据库和 RPC,但取消只表示调用者不再等待,不证明提交已回滚。
3. Outbox 把业务状态与“发送意图”放进一个事务
直接“先更新数据库,再发消息”会在两步之间崩溃;反过来则可能发出不存在的业务状态。Outbox 在业务数据库内建立事件表,让业务更新与待发送事件共同提交:
CREATE TABLE outbox_event (
event_id CHAR(36) PRIMARY KEY,
aggregate_type VARCHAR(40) NOT NULL,
aggregate_id VARCHAR(64) NOT NULL,
event_type VARCHAR(80) NOT NULL,
payload JSON NOT NULL,
state VARCHAR(16) NOT NULL DEFAULT 'pending',
attempts INTEGER NOT NULL DEFAULT 0,
available_at TIMESTAMP NOT NULL,
created_at TIMESTAMP NOT NULL,
sent_at TIMESTAMP NULL
);
CREATE INDEX outbox_pending_idx
ON outbox_event (state, available_at, created_at);
一次支付确认的事务只改变订单并插入事件。条件更新避免重复状态跃迁,event_id 在重试之间保持不变:
BEGIN;
UPDATE orders
SET status = 'paid', version = version + 1
WHERE order_id = ? AND status = 'pending' AND version = ?;
-- 应检查影响行数恰好为 1
INSERT INTO outbox_event
(event_id, aggregate_type, aggregate_id, event_type, payload,
state, available_at, created_at)
VALUES (?, 'order', ?, 'order.paid', ?, 'pending', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP);
COMMIT;
此时事件尚未到 broker,但发送意图不会丢。事务提交前进程退出,两项都不存在;提交后退出,两项都存在,Relay 之后可以继续发布。
4. Relay 的领取、发布和确认生命周期
Relay 周期性领取到期事件,发布后标为 sent。多个实例可通过 FOR UPDATE SKIP LOCKED 分摊工作;不支持该语法的数据库可用租约列和条件更新。领取批次要小,事务内只改变“归属/租约”,不要持锁等待 broker 网络。
BEGIN;
SELECT event_id, event_type, payload
FROM outbox_event
WHERE state = 'pending' AND available_at <= CURRENT_TIMESTAMP
ORDER BY created_at
LIMIT 100
FOR UPDATE SKIP LOCKED;
UPDATE outbox_event
SET state = 'publishing', attempts = attempts + 1
WHERE event_id IN (?, ?, ?);
COMMIT;
broker 确认后再更新 sent。若发布成功而更新 sent 前崩溃,事件会再次发布,这是不可消除的崩溃窗口,所以 Outbox 通常提供“至少一次”而非恰好一次。发布失败则按有上限的指数退避重置 pending/available_at;超过阈值进入 dead 状态并告警,不能无限快速重试。
Relay 的 goroutine 必须受应用 context 管理,关闭时停止领取、等待在途发布完成,并释放租约。批量大小、并发数和 broker producer 应有明确上限,否则数据库恢复时积压会瞬间压垮下游。
5. 消费者用去重记录保护本地副作用
消费者不能只依赖 broker 的 exactly-once 宣称,因为数据库更新通常不在 broker 事务内。把消费记录与业务变更放入消费者自己的本地事务:
BEGIN;
INSERT INTO consumed_event (consumer, event_id, consumed_at)
VALUES ('inventory', ?, CURRENT_TIMESTAMP)
ON CONFLICT (consumer, event_id) DO NOTHING;
-- 仅当插入影响 1 行时执行库存变更
UPDATE inventory
SET available = available - ?, version = version + 1
WHERE sku = ? AND available >= ?;
COMMIT;
去重键必须包含消费者身份,因为不同消费者都需要处理同一事件。记录保留期至少覆盖消息最大重投期、备份恢复窗口和人工重放窗口。消费成功后再 ack;事务失败则 nack/不确认。邮件等无法纳入数据库事务的副作用,应再写本服务的 Outbox,而不是提交后直接发送。
事件 schema 应带 event_id、类型、发生时间、聚合 ID 和版本。消费者必须容忍新增字段,并对无法识别的版本进入隔离队列,不能悄悄确认。不同聚合通常不保证全局顺序;确需顺序时按 aggregate key 分区,并用聚合版本拒绝倒退事件。
6. 用 Go 实现短本地事务和幂等命令
database/sql 的 *sql.DB 是连接池,事务会独占一条连接直到提交或回滚。请求 context 必须作为第一个参数贯穿;事务内不调用外部 HTTP。
func markPaid(ctx context.Context, db *sql.DB, orderID, eventID string, version int64, payload []byte) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin mark paid: %w", err)
}
defer tx.Rollback()
result, err := tx.ExecContext(ctx, updateOrderSQL, orderID, version)
if err != nil {
return fmt.Errorf("update order %q: %w", orderID, err)
}
changed, err := result.RowsAffected()
if err != nil {
return fmt.Errorf("read affected rows: %w", err)
}
if changed != 1 {
return ErrOrderStateChanged
}
if _, err := tx.ExecContext(ctx, insertOutboxSQL, eventID, orderID, payload); err != nil {
return fmt.Errorf("insert outbox event %q: %w", eventID, err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit mark paid: %w", err)
}
return nil
}
唯一键冲突应按具体 driver 的结构化错误码分类。若同一个 eventID 已存在,读取其关联订单和 payload hash:完全相同可视为幂等成功,不同则是键复用冲突。不要通过匹配错误字符串判断。Commit 返回网络错误时结果未知,应按业务键查询最终状态。
7. Saga 是显式的长期状态机
Saga 将“创建订单、扣款、预留库存、安排发货”拆成独立本地事务。编排式 Saga 由 coordinator 保存当前步骤并发命令;协同式 Saga 由事件驱动下一服务。编排更容易看清全局状态和人工恢复;协同耦合较松,但事件链长后难追踪。两者都没有自动回滚,补偿是新的业务动作。
pending -> charging -> reserving -> confirmed
| |
v v
charge_failed compensating
|
refunding + releasing
v
cancelled
每次跃迁用 WHERE saga_id=? AND state=? AND version=? 条件更新,保证并发 worker 只有一个获胜。命令带 saga_id/step/attempt 作为幂等键。补偿顺序通常与正向步骤相反,但“退款”可能失败或需要数日,因此状态机必须允许 compensating、manual_review,而不是假装一次函数返回就恢复原状。
8. Saga 补偿必须基于业务语义
扣款的补偿是退款,不是删除支付记录;库存释放要确认预留仍属于当前 Saga;优惠券恢复可能受有效期变化影响。不可逆动作,如已发送实物、第三方到账或已读通知,需要前置审批、延后执行或人工流程。
协调器处理一步时采用“读取状态、执行幂等远端命令、记录结果”。不能持有协调器数据库事务跨越 RPC。远端成功而结果记录失败时会重试同一命令,所以远端幂等是协议的一部分。错误应分为业务终止、可重试瞬时错误、未知结果和永久基础设施错误;只有明确瞬时错误自动退避重试。
取消请求也不是删除 Saga。用户取消会写入期望状态,worker 在安全点转入补偿;已在执行的远端操作仍通过查询和幂等键收敛。整个 Saga 有业务截止时间,但单次 RPC 使用更短 deadline,预留时间做状态落库。
9. TCC 的 Try、Confirm、Cancel 与三类异常
TCC 适合能预留资源的场景。Try 校验并冻结额度;Confirm 把冻结转成正式扣减;Cancel 释放冻结。三个接口都必须幂等,且状态由唯一分支键保护:
INSERT INTO reservation (branch_id, account_id, amount, state, expires_at)
VALUES (?, ?, ?, 'tried', ?)
ON CONFLICT (branch_id) DO NOTHING;
UPDATE reservation SET state = 'confirmed'
WHERE branch_id = ? AND state = 'tried';
UPDATE reservation SET state = 'cancelled'
WHERE branch_id = ? AND state = 'tried';
空回滚是 Cancel 先于 Try 到达:应插入 cancelled 屏障,使迟到 Try 不再预留。悬挂是 Try 因网络延迟在 Cancel 后才执行,也由屏障阻止。重复调用则返回当前终态。Confirm 与 Cancel 竞争时只能一个条件更新成功,另一个读取终态并按协议响应。
预留有过期时间,但不能仅靠 TTL 自动删除;协调器可能正在 Confirm。清理任务应根据全局事务状态决定释放,并保留审计记录。TCC 增加存储字段、可用额度计算和热点锁竞争,只有业务确实需要实时预留时才使用。
10. DTM 的职责和子事务屏障
DTM client v1.18.7 提供 Saga、TCC、事务消息等客户端 API,服务端负责全局事务记录、分支调度和重试。它减少协调样板,但不会替业务定义补偿、幂等键或隔离级别。版本应在 go.mod 固定,并在升级时跑协议兼容与故障注入测试:
go get github.com/dtm-labs/client@v1.18.7
go mod tidy
go test ./...
子事务屏障在业务库记录 gid、branch ID、operation 和 barrier ID,通过唯一约束处理重复、空回滚与悬挂。分支处理器必须把屏障记录和业务变更放在同一个本地事务。DTM 服务不可达时,新全局事务应快速失败或进入明确待处理状态;已受理事务依靠持久化状态恢复,不能只靠进程内 goroutine。
HTTP Saga 用 NewSaga 声明正向和补偿端点,Submit 只表示协调器受理或返回错误;分支仍需自行鉴权和幂等:
func submitOrderSaga(dtmServer, gid, service string, request OrderRequest) error {
saga := dtmcli.NewSaga(dtmServer, gid).
Add(service+"/charge", service+"/refund", request).
Add(service+"/reserve", service+"/release", request)
if err := saga.Submit(); err != nil {
return fmt.Errorf("submit saga %q: %w", gid, err)
}
return nil
}
接入时需固定 DTM server 与 client 的兼容组合,给回调地址做服务身份认证和网络限制。全局事务 ID 不授予权限;分支回调仍需验证租户、资源归属和操作范围,防止伪造 Confirm/Cancel。
11. 并发、连接池和背压是一套容量设计
Outbox Relay、Saga worker 和 TCC 回调都会竞争数据库连接。每实例 MaxOpenConns 乘实例数必须小于数据库总连接预算并留出管理余量。长事务、未关闭 Rows 和锁等待会提高 DB.Stats().WaitCount/WaitDuration,盲目扩池只会把更多并发推给数据库。
worker 并发应受 semaphore 或固定数量 goroutine 限制。领取 100 条不代表同时发 100 个请求;应按下游配额、单次延迟和超时预算测算。积压恢复时采用速率限制和随机抖动,避免所有实例同时扫描 pending 首页。热点 aggregate 可按键分片;同一键上的命令串行化仍需数据库版本条件兜底。
不要把 *sql.Tx 交给多个 goroutine 并发使用。事务占用单连接,driver 通常不支持并发语句;并发远端调用也不应包在本地事务中。先完成受控远端阶段,再用短事务记录结果。
12. 错误、取消、超时与安全边界
错误包装提供操作上下文并保留 %w,在边界用 errors.Is/As 或 driver 错误码分类;同一错误只记录或返回一次。日志记录 gid、event ID、branch、state、attempt、错误类别和耗时,不记录支付凭证、完整 payload、DSN 或用户隐私。
入口设置总 deadline,分支调用从剩余预算派生更短超时。达到 deadline 后停止当前等待并持久化 unknown/retryable 状态,不能直接执行相反操作。后台恢复使用服务生命周期 context,关闭时 cancel 并等待 worker;不能从请求 handler 启动无人管理的 goroutine。
Outbox payload 是不可信输入:消费者校验 schema、大小、枚举和租户;动态 SQL 标识符只能来自白名单,值始终参数化。管理端的重放、跳过和强制补偿属于高危写操作,需要鉴权、双重确认和不可变审计。回调接口使用 mTLS 或签名、时间窗和 nonce,并限制来源网络及请求体大小。
13. 迁移和发布必须兼容在途事务
分布式流程可能跨越多个版本。事件先新增可选字段并部署兼容消费者,再部署生产者;删除字段要等最大重放窗口过去。状态枚举新增时旧 worker 必须把未知状态留给新版处理,不能默认标成失败。
数据库迁移采用 expand/contract:先加可空列或新表、回填、双读验证,再切换写入,最后收紧约束。为 Outbox 建索引要评估在线 DDL 与写放大。修改 Saga 步骤时给流程定义加版本,新实例继续按旧定义处理旧 Saga,新请求才使用新定义。
灾备恢复也会重放旧事件。恢复数据库与 broker 到不同时间点时,依靠业务键和消费者去重收敛,并运行对账任务。不要清理仍可能被备份恢复引用的去重和屏障记录。
14. 诊断、对账与生产运行边界
至少监控:Outbox 最老 pending 年龄与数量、发布成功率和重试分布、消费延迟、去重命中、Saga 各状态停留时间、补偿失败、TCC 预留年龄、数据库池等待和锁等待。Trace 用 gid/event_id 关联,但业务状态以数据库审计记录为准,不能只依赖采样 Trace。
故障诊断按生命周期查:业务事务是否提交、Outbox 是否存在、Relay 是否领取、broker 是否确认、消费者本地事务是否提交、ack 是否返回。每一步都需要可查询证据。建立定时对账,以支付提供方流水、订单和账本的不变量找差异;修复通过幂等命令执行,禁止直接手改多张表。
生产边界包括积压上限、单事件大小、最大尝试次数、Saga 总时限和人工 SLA。下游大面积故障时应暂停领取或打开熔断,保留数据而不是持续制造超时。所谓最终一致必须给出“最终多久”和超限后的责任路径。
15. 测试状态机、SQL 契约和故障窗口
纯逻辑先用表驱动测试固定允许的状态跃迁、重复 Confirm/Cancel 和乱序事件。SQL 层用真实数据库容器测试唯一约束、隔离、SKIP LOCKED、死锁错误码和迁移;内存 fake 不能证明这些数据库语义。
func next(current, command string) (string, error) {
transitions := map[string]map[string]string{
"pending": {"charge": "charging", "cancel": "cancelled"},
"charging": {"charged": "reserving", "failed": "cancelled"},
"reserving": {"reserved": "confirmed", "failed": "compensating"},
"compensating": {"compensated": "cancelled"},
}
if target, ok := transitions[current][command]; ok {
return target, nil
}
return "", fmt.Errorf("transition %s with %s: %w", current, command, ErrInvalidTransition)
}
故障注入要覆盖:业务提交前退出、提交后发布前退出、发布后标记前退出、消费提交后 ack 前退出、RPC 成功响应丢失、两个 worker 同时领取、Confirm/Cancel 乱序、数据库主从切换及进程优雅关闭。断言不是“没有报错”,而是余额守恒、库存不为负、每个全局事务只有一个终态、重复执行不产生额外副作用。
性能测试同时测正常流量和积压恢复,观察数据库 CPU、索引扫描、连接池等待、broker ack 延迟及下游限流。优化前先确保状态可恢复;用跳过审计或删除去重记录换吞吐,会把暂时性能问题变成永久数据错误。
分布式事务的工程核心,是承认不确定结果并把它持久化:本地事务保证一个边界内原子,Outbox 保证发送意图不丢,幂等与屏障吸收重复和乱序,Saga/TCC 表达跨边界状态,诊断与对账处理自动机制之外的现实。只有这些生命周期都能被测试和观察,框架提供的“事务 API”才真正有意义。
系列导航与关联阅读
- 系列入口:Go 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go 服务韧性设计:超时、重试、限流、熔断与隔离舱
- 下一篇:Go GORM 完整指南:模型、查询、事务、关联与性能边界
- 延伸:Go database/sql 基础:连接池、事务、Context 与 NULL
- 延伸:Go 可靠消息统一设计:Outbox、幂等、重试、顺序与死信
- 延伸:Go RabbitMQ 实战:Exchange、Queue、确认、重试与死信
- 延伸:Go Kafka 完整指南:分区、消费者组、Offset 与 franz-go
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论