微服务分布式事务的核心挑战
单体架构中,数据库本地事务通过ACID保证数据一致性。微服务拆分后,一次业务操作可能涉及订单、库存、支付等多个独立数据库,本地事务无法跨服务生效。直接使用分布式锁或两阶段提交(2PC)会带来严重的性能损耗和可用性风险。Go微服务生态中,Saga模式和本地消息表是两种被广泛验证的分布式事务方案,各有明确的适用场景。
Saga模式:长事务的编排与协调
Saga模式将长事务拆分为一系列本地事务,每个本地事务提交后通过消息触发下一步操作。任意步骤失败时,执行已完成步骤的补偿事务回滚。Saga分为编排式(Choreography)和协调式(Orchestration)两种实现。
编排式Saga通过事件驱动,各服务监听事件并决定下一步操作。优点是耦合度低,缺点是业务流程分散在多个服务中,难以追踪完整链路。
协调式Saga由一个中心协调器控制流程走向。以下为Go语言实现的订单创建协调式Saga:
package saga
import (
"context"
"fmt"
)
type Step struct {
Name string
Execute func(ctx context.Context, order *Order) error
Compensate func(ctx context.Context, order *Order) error
}
type Saga struct {
Steps []Step
}
func (s *Saga) Run(ctx context.Context, order *Order) error {
completed := 0
for i, step := range s.Steps {
if err := step.Execute(ctx, order); err != nil {
// 执行失败,反向补偿已完成的步骤
for j := i - 1; j >= 0; j-- {
if compErr := s.Steps[j].Compensate(ctx, order); compErr != nil {
// 补偿失败记录日志,人工介入
fmt.Printf("compensate failed: step=%s, err=%v\n",
s.Steps[j].Name, compErr)
}
}
return fmt.Errorf("saga step %q failed: %w", step.Name, err)
}
completed = i + 1
order.ExecutedSteps = completed
}
return nil
}
// 订单创建Saga定义
func NewOrderSaga(
inventorySvc InventoryService,
paymentSvc PaymentService,
) *Saga {
return &Saga{
Steps: []Step{
{
Name: "reserve_inventory",
Execute: func(ctx context.Context, o *Order) error {
return inventorySvc.Reserve(ctx, o.ID, o.Items)
},
Compensate: func(ctx context.Context, o *Order) error {
return inventorySvc.Release(ctx, o.ID)
},
},
{
Name: "charge_payment",
Execute: func(ctx context.Context, o *Order) error {
return paymentSvc.Charge(ctx, o.ID, o.Total)
},
Compensate: func(ctx context.Context, o *Order) error {
return paymentSvc.Refund(ctx, o.ID)
},
},
},
}
}
本地消息表:最终一致性的可靠方案
本地消息表方案的核心思路是:在业务数据库中创建一张消息表,将业务操作和消息写入放在同一个本地事务中,通过后台任务轮询消息表并发送消息至消息中间件。发送成功的消息标记为已完成,失败则重试。这种方案保证业务操作和消息发送的原子性,无需引入额外的分布式事务协调组件。
以下为Go语言实现的本地消息表核心逻辑:
package outbox
import (
"context"
"database/sql"
"encoding/json"
"time"
)
type OutboxMessage struct {
ID int64
Topic string
Key string
Payload string
Status string // pending, sent, failed
Retry int
CreatedAt time.Time
}
type OutboxRepo struct {
db *sql.DB
}
// 在业务事务中写入消息,保证原子性
func (r *OutboxRepo) SaveWithTx(
ctx context.Context,
tx *sql.Tx,
topic, key string,
payload interface{},
) error {
data, err := json.Marshal(payload)
if err != nil {
return err
}
_, err = tx.ExecContext(ctx,
`INSERT INTO outbox_messages (topic, msg_key, payload, status, retry, created_at)
VALUES (?, ?, ?, 'pending', 0, NOW())`,
topic, key, string(data),
)
return err
}
// 后台轮询发送未完成消息
func (r *OutboxRepo) PollAndSend(
ctx context.Context,
sender func(topic, key, payload string) error,
) error {
rows, err := r.db.QueryContext(ctx,
`SELECT id, topic, msg_key, payload, retry
FROM outbox_messages
WHERE status = 'pending' AND retry < 5
ORDER BY created_at ASC
LIMIT 100`,
)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var m OutboxMessage
if err := rows.Scan(&m.ID, &m.Topic, &m.Key, &m.Payload, &m.Retry); err != nil {
continue
}
if err := sender(m.Topic, m.Key, m.Payload); err != nil {
// 发送失败,增加重试计数
r.db.ExecContext(ctx,
`UPDATE outbox_messages SET retry = retry + 1 WHERE id = ?`,
m.ID,
)
continue
}
// 发送成功,标记完成
r.db.ExecContext(ctx,
`UPDATE outbox_messages SET status = 'sent' WHERE id = ?`,
m.ID,
)
}
return nil
}
Saga与本地消息表的场景对比
一致性要求。Saga提供业务层面的回滚能力,适合需要即时感知失败并撤销的场景(如库存扣减后支付失败的释放)。本地消息表只保证最终一致性,适合对实时性要求不高的场景(如积分发放、通知推送)。
实现复杂度。Saga需要定义每个步骤的补偿逻辑,补偿链越长开发成本越高。本地消息表实现简单,但需要处理消息幂等性和重复消费问题。
性能与可用性。Saga的协调器是单点,需要实现高可用部署。本地消息表无单点依赖,但消息轮询有延迟。
可观测性。Saga协调器可追踪完整的执行与补偿链路。本地消息表的消息状态可通过数据库查询,但跨服务的完整链路追踪需要额外的分布式追踪支持。
混合方案与工程建议
在实际项目中,两种方案可以混合使用。核心业务链路(订单→库存→支付)采用Saga模式保证回滚能力,非核心链路(积分、通知、审计日志)采用本地消息表实现最终一致性。关键原则是:业务侧的可逆操作优先Saga,不可逆操作(如已发起的银行转账)优先本地消息表配合对账机制。无论选择哪种方案,消息幂等性、重试机制、死信队列处理都是必须实现的工程基线。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-wei-fu-wu-fen-bu-shi-shi-wu-shi-jian-saga-mo-shi-yu-ben/