#实战:多数据库协作的 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 旁路日志。
- 写路径 = 事务保证核心一致性 + 删缓存 + 异步日志;读路径 = 缓存优先 + 击穿/穿透/雪崩防护。
- 跨库采用最终一致,而非分布式强一致。
- 生产化离不开连接池、超时、幂等、可观测、降级与安全加固。