分布式事务怎么搞?Go后端新人也能看懂的最佳实践

Tech数据
2025-12-29 16:27
阅读 1450

大家好,我是你们的学长,一名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 对业务侵入大,新手容易出错,建议先掌握本地消息表。


六、学习建议与避坑指南

下一步学什么?

  1. 深入消息队列:学 RabbitMQ / Kafka 的 Exactly-Once 语义
  2. Seata 框架:阿里开源的分布式事务中间件,支持 AT/TCC/Saga 模式
  3. Saga 模式:长事务的另一种解法,适合业务流程长的场景

面试题挑战准备

  • “分布式事务有哪些方案?优缺点?”
  • “2PC 和 3PC 的区别?”
  • “如何保证消息不丢失、不重复?”
  • “最终一致性 vs 强一致性,如何选择?”

💬 我面试时就被问过:“如果本地消息表写成功,但服务崩溃了,消息还能发出去吗?”
答案是:能! 因为后台任务会定期扫描未发送的消息,这就是“可靠事件模式”的核心。

避坑提醒

  • ❌ 不要一上来就用 2PC —— 性能瓶颈大
  • ✅ 优先考虑“业务可接受最终一致性”的场景
  • ✅ 所有远程调用都要考虑超时、重试、熔断
  • ✅ 日志要打全!方便排查“到底哪一步失败了”

结语

分布式事务听起来吓人,但拆开来看,无非是用各种手段保证多步操作的可靠性。作为后端开发者,你不需要一开始就精通所有方案,但一定要理解为什么需要它,以及最常用的“本地消息表”是如何工作的

希望这篇技术分享能帮你迈出第一步。如果你觉得有用,欢迎关注我的博客,后续我会继续用 Go 带你攻克更多后端难题!

最后送大家一句话:“分布式系统的坑,都是踩出来的;但别人的坑,你可以绕着走。”

加油,未来的后端工程师!

评论 0

最热最新
暂无评论
Tech数据Lv.1
0
影响力
0
文章
0
粉丝