实战:多数据库协作的 Go 服务

本章把前文知识串起来,构建一个简化的订单服务:MySQL 负责交易数据、Redis 承担缓存、MongoDB 记录事件日志。重点不在业务复杂度,而在多数据库协作的工程模式

架构设计

                 ┌────────────────────────────────────┐
   客户端 ─────► │            Order Service            │
                 │   GET /orders/{id}  ──► 读路径      │
                 │   POST /orders      ──► 写路径      │
                 └───────┬──────────┬──────────┬───────┘
                         │          │          │
            ┌────────────▼──┐  ┌────▼─────┐  ┌─▼──────────┐
            │   MySQL       │  │  Redis   │  │  MongoDB   │
            │ 订单/库存     │  │ 订单缓存 │  │ 事件日志   │
            │ 强一致、事务  │  │ 低延迟   │  │ 半结构化   │
            └───────────────┘  └──────────┘  └────────────┘
数据库角色选择理由
MySQL订单、库存等核心交易数据事务、强一致、约束
Redis订单详情缓存、限流计数低延迟、高吞吐
MongoDB下单/支付事件日志半结构化、写入吞吐高、schema 灵活

表结构与数据模型

CREATE TABLE orders (
    id         BIGINT       PRIMARY KEY,
    user_id    BIGINT       NOT NULL,
    amount     DECIMAL(12,2) NOT NULL,
    status     VARCHAR(16)  NOT NULL,
    created_at DATETIME(3)  NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
    KEY idx_user_created (user_id, created_at)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

CREATE TABLE inventory (
    sku BIGINT PRIMARY KEY,
    qty INT    NOT NULL,
    CONSTRAINT chk_qty CHECK (qty >= 0)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
// 事件日志文档(半结构化,字段可随业务扩展)
{
  _id: ObjectId(),
  type: "order.created",          // 事件类型
  orderId: 1001,
  userId: 7,
  payload: { sku: 9001, qty: 2, amount: NumberDecimal("199.00") },
  at: new Date()
}

基础设施初始化

package main

import (
    "context"
    "os"
    "time"

    "github.com/redis/go-redis/v9"
    "go.mongodb.org/mongo-driver/mongo"
    "go.mongodb.org/mongo-driver/mongo/options"
    "gorm.io/driver/mysql"
    "gorm.io/gorm"
)

type Infra struct {
    DB     *gorm.DB
    Cache  *redis.Client
    Events *mongo.Collection
}

func NewInfra(ctx context.Context) (*Infra, error) {
    // MySQL(通过 DSN 环境变量注入)
    gdb, err := gorm.Open(mysql.Open(os.Getenv("MYSQL_DSN")), &gorm.Config{})
    if err != nil {
        return nil, err
    }
    sqlDB, _ := gdb.DB()
    sqlDB.SetMaxOpenConns(25)
    sqlDB.SetMaxIdleConns(25)
    sqlDB.SetConnMaxLifetime(5 * time.Minute)

    // Redis
    cache := redis.NewClient(&redis.Options{Addr: os.Getenv("REDIS_ADDR"), PoolSize: 20})
    if err := cache.Ping(ctx).Err(); err != nil {
        return nil, err
    }

    // MongoDB
    client, err := mongo.Connect(ctx, options.Client().ApplyURI(os.Getenv("MONGO_URI")))
    if err != nil {
        return nil, err
    }
    if err := client.Ping(ctx, nil); err != nil {
        return nil, err
    }

    return &Infra{
        DB:     gdb,
        Cache:  cache,
        Events: client.Database("shop").Collection("events"),
    }, nil
}
Tip

生产环境的连接参数(DSN、地址、口令)全部通过环境变量或配置中心注入,代码里不出现任何明文凭据。三类客户端都应在启动时 Ping 验证并配置连接池与生命周期。

订单模型

package main

import "time"

type Order struct {
    ID        int64     `gorm:"primaryKey" json:"id"`
    UserID    int64     `gorm:"index:idx_user_created" json:"user_id"`
    Amount    float64   `json:"amount"`
    Status    string    `gorm:"size:16" json:"status"`
    CreatedAt time.Time `json:"created_at"`
}

type Event struct {
    Type    string         `bson:"type"`
    OrderID int64          `bson:"orderId"`
    UserID  int64          `bson:"userId"`
    Payload map[string]any `bson:"payload"`
    At      time.Time      `bson:"at"`
}

写路径:事务 + 缓存失效 + 日志

下单要保证「扣库存 + 建订单」的原子性,写成功后删除缓存,并异步写事件日志。

package main

import (
    "context"
    "fmt"
    "time"

    "gorm.io/gorm"
)

type CreateOrderReq struct {
    UserID int64
    SKU    int64
    Qty    int
    Amount float64
}

func (s *OrderService) CreateOrder(ctx context.Context, req CreateOrderReq) (*Order, error) {
    var order Order

    // 1. 事务:原子扣减库存 + 创建订单
    err := s.infra.DB.Transaction(func(tx *gorm.DB) error {
        // 条件更新,防止超卖:影响行数为 0 表示库存不足
        res := tx.Exec(
            "UPDATE inventory SET qty = qty - ? WHERE sku = ? AND qty >= ?",
            req.Qty, req.SKU, req.Qty,
        )
        if res.Error != nil {
            return res.Error
        }
        if res.RowsAffected == 0 {
            return ErrInsufficientStock
        }

        order = Order{
            ID:     s.nextID(ctx), // 雪花 ID / 号段,趋势递增
            UserID: req.UserID,
            Amount: req.Amount,
            Status: "created",
            CreatedAt: time.Now(),
        }
        return tx.Create(&order).Error
    })
    if err != nil {
        return nil, err
    }

    // 2. 删除缓存(Cache Aside:写库后删缓存,下次读回填)
    cacheKey := fmt.Sprintf("order:%d", order.ID)
    s.infra.Cache.Del(ctx, cacheKey)

    // 3. 异步写事件日志(不阻塞主流程,失败可重试)
    s.emitEventAsync(ctx, Event{
        Type:    "order.created",
        OrderID: order.ID,
        UserID:  order.UserID,
        Payload: map[string]any{"sku": req.SKU, "qty": req.Qty, "amount": req.Amount},
        At:      time.Now(),
    })

    return &order, nil
}

func (s *OrderService) emitEventAsync(ctx context.Context, e Event) {
    go func() {
        c, cancel := context.WithTimeout(context.Background(), 3*time.Second)
        defer cancel()
        if _, err := s.infra.Events.InsertOne(c, e); err != nil {
            slog.Error("写入事件日志失败", "err", err, "orderId", e.OrderID)
            // 生产:写入重试队列 / 死信,保证最终写入
        }
    }()
}

读路径:缓存优先 + 击穿保护

package main

import (
    "context"
    "encoding/json"
    "errors"
    "fmt"
    "time"

    "github.com/redis/go-redis/v9"
    "gorm.io/gorm"
)

var ErrNotFound = errors.New("order not found")

func (s *OrderService) GetOrder(ctx context.Context, id int64) (*Order, error) {
    key := fmt.Sprintf("order:%d", id)

    // 1. 命中缓存直接返回
    if data, err := s.infra.Cache.Get(ctx, key).Bytes(); err == nil {
        var o Order
        if json.Unmarshal(data, &o) == nil {
            return &o, nil
        }
    } else if err != redis.Nil {
        slog.Warn("读取缓存失败,降级查库", "err", err)
    }

    // 2. 互斥锁防止缓存击穿:仅一个请求回源
    lockKey := key + ":lock"
    locked, _ := s.infra.Cache.SetNX(ctx, lockKey, "1", 5*time.Second).Result()
    if !locked {
        // 未抢到锁,短暂等待后重试读缓存
        time.Sleep(30 * time.Millisecond)
        if data, err := s.infra.Cache.Get(ctx, key).Bytes(); err == nil {
            var o Order
            if json.Unmarshal(data, &o) == nil {
                return &o, nil
            }
        }
    } else {
        defer s.infra.Cache.Del(ctx, lockKey)
    }

    // 3. 查库
    var o Order
    err := s.infra.DB.First(&o, "id = ?", id).Error
    if errors.Is(err, gorm.ErrRecordNotFound) {
        // 缓存空值防穿透,短过期
        s.infra.Cache.Set(ctx, key, []byte{}, 60*time.Second)
        return nil, ErrNotFound
    }
    if err != nil {
        return nil, err
    }

    // 4. 回填缓存,过期时间加随机抖动防雪崩
    data, _ := json.Marshal(o)
    ttl := 10*time.Minute + time.Duration(rand.Intn(120))*time.Second
    s.infra.Cache.Set(ctx, key, data, ttl)

    return &o, nil
}

HTTP 层

使用 Go 1.22+ 增强的 net/http 路由,无需第三方框架。

package main

import (
    "context"
    "encoding/json"
    "errors"
    "log/slog"
    "net/http"
)

func main() {
    infra, err := NewInfra(context.Background())
    if err != nil {
        slog.Error("初始化失败", "err", err)
        return
    }
    svc := &OrderService{infra: infra}

    mux := http.NewServeMux()
    mux.HandleFunc("POST /orders", func(w http.ResponseWriter, r *http.Request) {
        var req CreateOrderReq
        if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
            http.Error(w, "invalid body", http.StatusBadRequest)
            return
        }
        order, err := svc.CreateOrder(r.Context(), req)
        if errors.Is(err, ErrInsufficientStock) {
            http.Error(w, "库存不足", http.StatusConflict)
            return
        }
        if err != nil {
            slog.Error("下单失败", "err", err)
            http.Error(w, "internal error", http.StatusInternalServerError)
            return
        }
        writeJSON(w, http.StatusCreated, order)
    })

    mux.HandleFunc("GET /orders/{id}", func(w http.ResponseWriter, r *http.Request) {
        id, err := parseID(r.PathValue("id"))
        if err != nil {
            http.Error(w, "invalid id", http.StatusBadRequest)
            return
        }
        order, err := svc.GetOrder(r.Context(), id)
        if errors.Is(err, ErrNotFound) {
            http.Error(w, "not found", http.StatusNotFound)
            return
        }
        if err != nil {
            http.Error(w, "internal error", http.StatusInternalServerError)
            return
        }
        writeJSON(w, http.StatusOK, order)
    })

    slog.Info("listening on :8080")
    _ = http.ListenAndServe(":8080", mux)
}

一致性策略回顾

场景策略
下单扣库存MySQL 事务 + 条件更新防超卖
缓存与库一致写库后删缓存(Cache Aside)+ 过期兜底
缓存穿透空值缓存 + 短过期
缓存击穿分布式互斥锁,仅一个请求回源
缓存雪崩过期时间随机抖动 + Redis 高可用
事件日志异步写入 + 失败重试,不阻塞主链路
跨库一致性不追求强一致;核心交易走 MySQL,日志走最终一致
Warning

不要试图用「分布式事务」把三个库绑在一起。跨库强一致代价极高且脆弱。正确做法是:核心交易数据只放一个权威库(这里是 MySQL),其他库承担缓存与旁路日志,通过「事务内只改权威库 + 事务后异步同步」实现最终一致

生产加固清单

  • 连接池:三类客户端都配置上限、空闲数与生命周期。
  • 超时:所有操作带 context 超时;HTTP 层加读写超时。
  • 幂等:下单接口支持幂等键(Idempotency Key),避免重复下单。
  • 可观测:结构化日志、指标(QPS/延迟/缓存命中率)、链路追踪。
  • 降级:Redis/MongoDB 不可用时,读降级到 MySQL、日志降级到本地队列。
  • 安全:最小权限账号、TLS、参数化查询、敏感字段脱敏。

小结

  • 多数据库各司其职:MySQL 权威存储、Redis 缓存、MongoDB 旁路日志。
  • 写路径 = 事务保证核心一致性 + 删缓存 + 异步日志;读路径 = 缓存优先 + 击穿/穿透/雪崩防护。
  • 跨库采用最终一致,而非分布式强一致。
  • 生产化离不开连接池、超时、幂等、可观测、降级与安全加固。