Go微服务分布式事务实践:Saga模式与本地消息表方案详解

微服务分布式事务的核心挑战

单体架构中,数据库本地事务通过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/

(0)
小编小编
上一篇 18小时前
下一篇 18小时前

相关推荐