返回首页

Golang工程化:5.6 缓存、消息队列、事务一致性设计

在分布式系统里,性能和一致性几乎总是在拉扯:你想让系统更快,就会引入缓存;你想解耦和削峰,就会引入消息队列;你想把跨服务操作做成一个“整体”,就会碰到分布式事务。

问题也恰恰从这里开始:缓存可能过期不及时,消息可能重复投递,事务可能只能做到最终一致而不是强一致。很多线上故障,并不是因为某一种技术不会用,而是因为把这些技术放在一起后,没有把边界条件和失败场景设计清楚。

本文围绕「缓存、消息队列、事务一致性设计」展开,重点讲清楚以下几个问题:

  • 什么时候应该用缓存,什么时候不该用缓存
  • 常见缓存一致性策略如何选型
  • 缓存击穿、穿透、雪崩分别是什么,怎么处理
  • 消息队列如何帮助系统实现最终一致性
  • 2PC、Saga、TCC 的核心差异与适用场景
  • 补偿机制和幂等设计为什么是工程落地的关键

一、什么时候用缓存,什么时候别用

缓存的本质,是用“空间换时间”,把高频访问数据放到更快的存储层里,以减少对数据库或下游服务的直接访问。

1.1 适合使用缓存的场景

以下几类场景通常非常适合引入缓存:

  • 读多写少:例如商品详情、文章内容、配置信息、用户基础资料
  • 热点明显:少量数据被大量访问,例如秒杀商品、首页推荐位、热门榜单
  • 计算代价高:结果需要复杂聚合、排序、远程调用,重复计算成本高
  • 允许短暂不一致:例如展示类数据允许几秒延迟
  • 下游系统较脆弱:用缓存吸收流量,降低数据库和 RPC 服务压力

1.2 不适合使用缓存的场景

缓存不是越多越好。以下情况要谨慎,甚至尽量不用:

  • 强一致要求极高:例如余额、库存最终扣减值、核心账务流水
  • 数据变化频繁且读取收益不高:刚写进去马上又失效,缓存命中率很低
  • 数据基数极大但访问离散:缓存空间浪费严重,冷热不均明显
  • 错误代价极高:一旦缓存脏读可能导致资损、超卖、错误决策
  • 链路已经非常复杂:再引入缓存会提高系统理解和排障难度

1.3 一个实用判断标准

可以用下面这张表来快速判断。

维度 适合缓存 不适合缓存
读写比 读多写少 写多读少
一致性要求 可接受短暂延迟 必须强一致
数据访问模式 热点集中 随机离散
数据构建成本
错误容忍度 展示类、统计类 资金类、交易类

一句话概括:缓存应该优先服务“高频读取、可容忍短暂不一致”的场景,而不是替代数据库承担真相存储。

二、缓存一致性策略:Cache-Aside、Write-Through、Write-Behind

缓存一致性不是指“缓存和数据库永远完全一样”,而是指:在系统可接受的业务语义下,让缓存和数据库尽量保持正确、可控、可恢复的状态。

2.1 Cache-Aside

Cache-Aside 是最常见的策略,也叫“旁路缓存”。

读取流程:

  1. 先查缓存
  2. 缓存命中则直接返回
  3. 缓存未命中则查数据库
  4. 把数据库结果写入缓存
  5. 返回结果

写入流程通常是:

  1. 先更新数据库
  2. 再删除缓存

注意,这里一般是删缓存而不是“更新缓存”。原因很简单:更新缓存容易带入更多并发时序问题,而删除缓存更稳妥。后续读请求再触发回填即可。

优点:

  • 实现简单,工程成本低
  • 适合绝大多数读多写少场景
  • 容易和现有数据库系统配合

缺点:

  • 不是强一致
  • 在高并发更新场景下可能出现短暂脏读
  • 需要处理缓存击穿、回源放大等问题

2.2 Write-Through

Write-Through 的思路是:写请求先进入缓存,再由缓存同步写数据库。 对业务方来说,缓存是主要写入口。

优点:

  • 读写路径更统一
  • 缓存中的数据通常更“新”
  • 对某些封装良好的缓存层,业务代码更简单

缺点:

  • 写延迟更高,因为要同步落库
  • 系统设计复杂度上升
  • 如果缓存层能力不足,容易成为瓶颈

2.3 Write-Behind

Write-Behind 也叫 Write-Back。业务先写缓存,数据库异步落盘。

优点:

  • 写性能非常高
  • 适合日志聚合、统计累加、批量落库等场景
  • 能显著降低数据库写压力

缺点:

  • 数据丢失风险更高
  • 故障恢复复杂
  • 一致性最弱,不适合核心交易数据

2.4 三种策略对比

策略 读路径 写路径 一致性 实现复杂度 典型场景
Cache-Aside 先缓存,后数据库 先数据库,后删缓存 中等 商品详情、用户资料
Write-Through 读缓存 先缓存,再同步落库 较高 封装式缓存服务
Write-Behind 读缓存 先缓存,异步落库 较低 计数、日志、批量写

2.5 可运行示例:三种缓存写法对比

下面的代码是一个完整可运行的 Go 示例。它使用内存 map 模拟数据库和缓存,演示 Cache-Aside、Write-Through、Write-Behind 三种策略的基本行为。

package main

import (
	"fmt"
	"sync"
	"time"
)

type User struct {
	ID   int
	Name string
	Age  int
}

type DB struct {
	mu   sync.RWMutex
	data map[int]User
}

func NewDB() *DB {
	return &DB{data: make(map[int]User)}
}

func (db *DB) Get(id int) (User, bool) {
	db.mu.RLock()
	defer db.mu.RUnlock()
	u, ok := db.data[id]
	return u, ok
}

func (db *DB) Set(u User) {
	db.mu.Lock()
	defer db.mu.Unlock()
	db.data[u.ID] = u
}

type CacheItem struct {
	Value    User
	ExpireAt time.Time
}

type Cache struct {
	mu   sync.RWMutex
	data map[int]CacheItem
}

func NewCache() *Cache {
	return &Cache{data: make(map[int]CacheItem)}
}

func (c *Cache) Get(id int) (User, bool) {
	c.mu.RLock()
	item, ok := c.data[id]
	c.mu.RUnlock()
	if !ok {
		return User{}, false
	}
	if !item.ExpireAt.IsZero() && time.Now().After(item.ExpireAt) {
		c.Delete(id)
		return User{}, false
	}
	return item.Value, true
}

func (c *Cache) Set(id int, u User, ttl time.Duration) {
	c.mu.Lock()
	defer c.mu.Unlock()
	c.data[id] = CacheItem{Value: u, ExpireAt: time.Now().Add(ttl)}
}

func (c *Cache) Delete(id int) {
	c.mu.Lock()
	defer c.mu.Unlock()
	delete(c.data, id)
}

type CacheAsideRepo struct {
	db    *DB
	cache *Cache
	ttl   time.Duration
}

func (r *CacheAsideRepo) GetUser(id int) (User, bool) {
	if u, ok := r.cache.Get(id); ok {
		fmt.Println("[Cache-Aside] hit cache")
		return u, true
	}
	fmt.Println("[Cache-Aside] miss cache, load from db")
	u, ok := r.db.Get(id)
	if !ok {
		return User{}, false
	}
	r.cache.Set(id, u, r.ttl)
	return u, true
}

func (r *CacheAsideRepo) UpdateUser(u User) {
	r.db.Set(u)
	r.cache.Delete(u.ID)
	fmt.Println("[Cache-Aside] db updated, cache deleted")
}

type WriteThroughRepo struct {
	db    *DB
	cache *Cache
	ttl   time.Duration
}

func (r *WriteThroughRepo) UpdateUser(u User) {
	r.cache.Set(u.ID, u, r.ttl)
	r.db.Set(u)
	fmt.Println("[Write-Through] cache updated, db updated synchronously")
}

func (r *WriteThroughRepo) GetUser(id int) (User, bool) {
	if u, ok := r.cache.Get(id); ok {
		fmt.Println("[Write-Through] hit cache")
		return u, true
	}
	fmt.Println("[Write-Through] miss cache, load from db")
	u, ok := r.db.Get(id)
	if !ok {
		return User{}, false
	}
	r.cache.Set(id, u, r.ttl)
	return u, true
}

type WriteBehindRepo struct {
	db      *DB
	cache   *Cache
	ttl     time.Duration
	flushCh chan User
	wg      sync.WaitGroup
}

func NewWriteBehindRepo(db *DB, cache *Cache, ttl time.Duration) *WriteBehindRepo {
	r := &WriteBehindRepo{
		db:      db,
		cache:   cache,
		ttl:     ttl,
		flushCh: make(chan User, 16),
	}
	r.wg.Add(1)
	go func() {
		defer r.wg.Done()
		for u := range r.flushCh {
			time.Sleep(100 * time.Millisecond)
			r.db.Set(u)
			fmt.Printf("[Write-Behind] async flush to db: %+v\n", u)
		}
	}()
	return r
}

func (r *WriteBehindRepo) UpdateUser(u User) {
	r.cache.Set(u.ID, u, r.ttl)
	r.flushCh <- u
	fmt.Println("[Write-Behind] cache updated, db flush scheduled")
}

func (r *WriteBehindRepo) GetUser(id int) (User, bool) {
	if u, ok := r.cache.Get(id); ok {
		fmt.Println("[Write-Behind] hit cache")
		return u, true
	}
	fmt.Println("[Write-Behind] miss cache, load from db")
	u, ok := r.db.Get(id)
	if !ok {
		return User{}, false
	}
	r.cache.Set(id, u, r.ttl)
	return u, true
}

func (r *WriteBehindRepo) Close() {
	close(r.flushCh)
	r.wg.Wait()
}

func main() {
	db := NewDB()
	db.Set(User{ID: 1, Name: "Alice", Age: 20})

	fmt.Println("===== Cache-Aside =====")
	ca := &CacheAsideRepo{db: db, cache: NewCache(), ttl: 3 * time.Second}
	u, _ := ca.GetUser(1)
	fmt.Println("read:", u)
	ca.UpdateUser(User{ID: 1, Name: "Alice", Age: 21})
	u, _ = ca.GetUser(1)
	fmt.Println("read after update:", u)

	fmt.Println("\n===== Write-Through =====")
	wt := &WriteThroughRepo{db: db, cache: NewCache(), ttl: 3 * time.Second}
	wt.UpdateUser(User{ID: 1, Name: "Alice", Age: 22})
	u, _ = wt.GetUser(1)
	fmt.Println("read:", u)

	fmt.Println("\n===== Write-Behind =====")
	wb := NewWriteBehindRepo(db, NewCache(), 3*time.Second)
	defer wb.Close()
	wb.UpdateUser(User{ID: 1, Name: "Alice", Age: 23})
	u, _ = wb.GetUser(1)
	fmt.Println("read immediately from cache:", u)
	time.Sleep(300 * time.Millisecond)
	u, _ = db.Get(1)
	fmt.Println("db after async flush:", u)
}

你可以把这段代码保存为 main.go 后执行:

go run main.go

从输出里你会直观看到三种策略的行为差异。

三、缓存击穿、穿透、雪崩及解决方案

这三个概念经常一起出现,但问题本质不同。

3.1 缓存击穿

定义: 某个热点 key 失效瞬间,大量请求同时回源数据库。

典型场景:

  • 热门商品详情突然过期
  • 热门配置项同一时刻失效
  • 秒杀活动的热点库存 key 被大量访问

常见解决方案:

  • 互斥锁 / singleflight:同一时刻只允许一个请求回源
  • 热点数据永不过期:用后台异步刷新代替被动失效
  • 提前续期:在过期前主动刷新热点 key

3.2 缓存穿透

定义: 请求的数据在缓存和数据库里都不存在,导致请求每次都穿透到数据库。

典型场景:

  • 恶意请求不存在的用户 ID
  • 参数非法但没有前置校验
  • 机器人批量探测系统边界

常见解决方案:

  • 参数校验:非法 ID 直接拦截
  • 空值缓存:数据库查不到时,缓存一个短 TTL 的空对象
  • 布隆过滤器:提前拦截明显不存在的数据

3.3 缓存雪崩

定义: 大量 key 在同一时间集中过期,导致回源流量瞬间压垮数据库或下游服务。

典型场景:

  • 批量预热时使用了相同 TTL
  • Redis 故障导致缓存整体不可用
  • 活动开始前后大量 key 集中失效

常见解决方案:

  • TTL 增加随机抖动,避免同一时刻过期
  • 多级缓存:本地缓存 + Redis
  • 限流、降级、熔断,保护数据库
  • 缓存集群高可用,避免单点故障

3.4 三者对比

问题 核心特征 风险点 典型解法
击穿 热点 key 失效 瞬时并发回源 singleflight、互斥锁、热点不过期
穿透 请求不存在的数据 无意义打 DB 空值缓存、布隆过滤器、参数校验
雪崩 大量 key 同时失效 整体流量冲垮下游 TTL 抖动、多级缓存、限流降级

3.5 可运行示例:击穿、穿透、雪崩保护

下面这个示例用标准库实现了一个带保护策略的缓存访问器,包含:

  • 非法参数拦截
  • 空值缓存,防止穿透
  • 单 key 互斥回源,防止击穿
  • TTL 抖动,降低雪崩风险
package main

import (
	"fmt"
	"math/rand"
	"sync"
	"time"
)

type Product struct {
	ID   int
	Name string
}

type cacheValue struct {
	product  Product
	exists   bool
	expireAt time.Time
}

type ProductService struct {
	db      map[int]Product
	cache   map[int]cacheValue
	mu      sync.RWMutex
	keyLock sync.Map
}

func NewProductService() *ProductService {
	return &ProductService{
		db: map[int]Product{
			1: {ID: 1, Name: "Keyboard"},
			2: {ID: 2, Name: "Mouse"},
		},
		cache: make(map[int]cacheValue),
	}
}

func (s *ProductService) getKeyLock(id int) *sync.Mutex {
	v, _ := s.keyLock.LoadOrStore(id, &sync.Mutex{})
	return v.(*sync.Mutex)
}

func randomTTL(base time.Duration) time.Duration {
	jitter := time.Duration(rand.Intn(500)) * time.Millisecond
	return base + jitter
}

func (s *ProductService) loadFromDB(id int) (Product, bool) {
	time.Sleep(50 * time.Millisecond)
	p, ok := s.db[id]
	return p, ok
}

func (s *ProductService) GetProduct(id int) (Product, bool) {
	if id <= 0 {
		fmt.Println("invalid id blocked")
		return Product{}, false
	}

	now := time.Now()
	s.mu.RLock()
	if v, ok := s.cache[id]; ok && now.Before(v.expireAt) {
		s.mu.RUnlock()
		if !v.exists {
			fmt.Println("negative cache hit")
			return Product{}, false
		}
		fmt.Println("cache hit")
		return v.product, true
	}
	s.mu.RUnlock()

	lock := s.getKeyLock(id)
	lock.Lock()
	defer lock.Unlock()

	s.mu.RLock()
	if v, ok := s.cache[id]; ok && now.Before(v.expireAt) {
		s.mu.RUnlock()
		if !v.exists {
			fmt.Println("negative cache hit after lock")
			return Product{}, false
		}
		fmt.Println("cache hit after lock")
		return v.product, true
	}
	s.mu.RUnlock()

	fmt.Println("cache miss, load from db")
	p, ok := s.loadFromDB(id)

	s.mu.Lock()
	defer s.mu.Unlock()
	if !ok {
		s.cache[id] = cacheValue{exists: false, expireAt: time.Now().Add(2 * time.Second)}
		return Product{}, false
	}
	s.cache[id] = cacheValue{product: p, exists: true, expireAt: time.Now().Add(randomTTL(3 * time.Second))}
	return p, true
}

func main() {
	rand.Seed(time.Now().UnixNano())
	svc := NewProductService()

	fmt.Println("===== cache penetration =====")
	_, ok := svc.GetProduct(-1)
	fmt.Println("result:", ok)
	_, ok = svc.GetProduct(999)
	fmt.Println("result:", ok)
	_, ok = svc.GetProduct(999)
	fmt.Println("result:", ok)

	fmt.Println("\n===== cache breakdown protection =====")
	var wg sync.WaitGroup
	for i := 0; i < 5; i++ {
		wg.Add(1)
		go func(i int) {
			defer wg.Done()
			p, ok := svc.GetProduct(1)
			fmt.Printf("goroutine=%d ok=%v product=%+v\n", i, ok, p)
		}(i)
	}
	wg.Wait()

	fmt.Println("\n===== cache avalanche mitigation hint =====")
	fmt.Println("TTL uses random jitter, so keys are less likely to expire at the same moment")
}

执行方式同样是:

go run main.go

这个示例虽然是简化版,但它已经体现了线上治理的几个关键思路:参数前置校验、空值缓存、回源锁、TTL 打散。

四、消息队列在一致性中的作用:事务消息、最终一致性

当业务跨多个服务时,本地数据库事务已经不够用了。

比如一个下单流程:

  1. 订单服务创建订单
  2. 库存服务扣减库存
  3. 营销服务发放优惠权益
  4. 通知服务发送消息

如果把这些动作放在一个调用链里同步完成,会遇到几个问题:

  • 调用链过长,延迟高
  • 任一服务失败会拖垮整体成功率
  • 很难做到跨服务强一致
  • 下游峰值会直接传导给上游

消息队列的价值,不只是“异步”,更重要的是:通过可靠投递 + 可重试 + 幂等消费,把跨服务事务问题转化为最终一致性问题。

4.1 事务消息是什么

事务消息的核心目标是:确保本地事务成功时,消息一定能发出去;本地事务失败时,消息一定不能错误发送。

常见实现有两类:

  1. MQ 原生事务消息

    • 由消息中间件提供半消息、事务回查等能力
    • 典型代表是部分支持事务消息的 MQ 产品
  2. Transactional Outbox(事务外盒)模式

    • 在本地事务里同时写业务表和 outbox 表
    • 后台任务把 outbox 表中的消息异步投递到 MQ
    • 这是非常常见、落地性强的做法

对工程实践来说,Outbox 往往比“依赖某个特定 MQ 的高级事务特性”更通用。

4.2 最终一致性是怎么实现的

最终一致性并不是“系统短时间内可以乱”,而是:

  • 状态变化有明确的传播路径
  • 中间失败可以重试
  • 重试不会导致副作用重复执行
  • 必要时可以补偿,把系统拉回可接受状态

所以,一个完整的最终一致性方案通常包括:

  • 可靠消息投递
  • 消费端幂等
  • 失败重试
  • 死信处理
  • 补偿逻辑
  • 状态机驱动而不是布尔值乱飞

4.3 消息队列不是“自动一致性”

很多人误以为“上了 MQ 就自动一致了”,这是很危险的。

消息队列只能保证“消息传递”这件事更可靠,并不能替你解决:

  • 消费重复执行
  • 下游部分成功、部分失败
  • 业务补偿顺序错误
  • 乱序消息覆盖新状态
  • 长时间重试造成脏状态

因此,MQ 是一致性的基础设施,不是一致性的最终答案。

4.4 可运行示例:Outbox + MQ + 幂等消费 + 补偿

下面这段代码是一个完整可运行的 Go 示例,模拟一个订单系统:

  • 创建订单时,本地事务内写入订单和 outbox 事件
  • 发布器异步把 outbox 推给“消息队列”
  • 库存消费者幂等处理订单创建事件
  • 库存不足时触发取消订单
  • 模拟重复投递,验证幂等性
package main

import (
	"fmt"
	"sync"
	"time"
)

type OrderStatus string

const (
	OrderPending   OrderStatus = "PENDING"
	OrderConfirmed OrderStatus = "CONFIRMED"
	OrderCanceled  OrderStatus = "CANCELED"
)

type Order struct {
	OrderID string
	SKU     string
	Qty     int
	Status  OrderStatus
}

type Event struct {
	EventID  string
	Type     string
	OrderID  string
	SKU      string
	Qty      int
	Sent     bool
	CreateAt time.Time
}

type OrderDB struct {
	mu     sync.Mutex
	orders map[string]Order
	outbox []Event
}

func NewOrderDB() *OrderDB {
	return &OrderDB{orders: make(map[string]Order)}
}

func (db *OrderDB) CreateOrderWithOutbox(order Order, event Event) {
	db.mu.Lock()
	defer db.mu.Unlock()
	db.orders[order.OrderID] = order
	db.outbox = append(db.outbox, event)
	fmt.Println("[OrderDB] local tx committed: order + outbox")
}

func (db *OrderDB) MarkOutboxSent(eventID string) {
	db.mu.Lock()
	defer db.mu.Unlock()
	for i := range db.outbox {
		if db.outbox[i].EventID == eventID {
			db.outbox[i].Sent = true
			return
		}
	}
}

func (db *OrderDB) PullUnsentEvents() []Event {
	db.mu.Lock()
	defer db.mu.Unlock()
	var events []Event
	for _, e := range db.outbox {
		if !e.Sent {
			events = append(events, e)
		}
	}
	return events
}

func (db *OrderDB) UpdateOrderStatus(orderID string, status OrderStatus) {
	db.mu.Lock()
	defer db.mu.Unlock()
	order := db.orders[orderID]
	order.Status = status
	db.orders[orderID] = order
}

func (db *OrderDB) GetOrder(orderID string) Order {
	db.mu.Lock()
	defer db.mu.Unlock()
	return db.orders[orderID]
}

type InventoryService struct {
	mu             sync.Mutex
	stock          map[string]int
	processedEvent map[string]bool
}

func NewInventoryService() *InventoryService {
	return &InventoryService{
		stock:          map[string]int{"sku-1": 2},
		processedEvent: make(map[string]bool),
	}
}

func (s *InventoryService) HandleOrderCreated(e Event) string {
	s.mu.Lock()
	defer s.mu.Unlock()

	if s.processedEvent[e.EventID] {
		fmt.Println("[Inventory] duplicate event ignored:", e.EventID)
		if s.stock[e.SKU] >= 0 {
			return "DUPLICATE"
		}
	}

	s.processedEvent[e.EventID] = true

	stock := s.stock[e.SKU]
	if stock < e.Qty {
		fmt.Println("[Inventory] reserve failed, insufficient stock")
		return "RESERVE_FAILED"
	}

	s.stock[e.SKU] -= e.Qty
	fmt.Println("[Inventory] reserve success, left stock:", s.stock[e.SKU])
	return "RESERVE_OK"
}

func (s *InventoryService) Release(orderID, sku string, qty int) {
	s.mu.Lock()
	defer s.mu.Unlock()
	s.stock[sku] += qty
	fmt.Printf("[Inventory] compensation release for order=%s, stock=%d\n", orderID, s.stock[sku])
}

func main() {
	orderDB := NewOrderDB()
	inventory := NewInventoryService()
	mq := make(chan Event, 16)

	order := Order{OrderID: "o-1001", SKU: "sku-1", Qty: 1, Status: OrderPending}
	event := Event{
		EventID:  "evt-1",
		Type:     "OrderCreated",
		OrderID:  order.OrderID,
		SKU:      order.SKU,
		Qty:      order.Qty,
		CreateAt: time.Now(),
	}

	orderDB.CreateOrderWithOutbox(order, event)

	go func() {
		for _, e := range orderDB.PullUnsentEvents() {
			mq <- e
			mq <- e // 模拟重复投递
			orderDB.MarkOutboxSent(e.EventID)
			fmt.Println("[Publisher] event published to mq:", e.EventID)
		}
	}()

	done := make(chan struct{})
	go func() {
		for i := 0; i < 2; i++ {
			e := <-mq
			result := inventory.HandleOrderCreated(e)
			switch result {
			case "RESERVE_OK":
				orderDB.UpdateOrderStatus(e.OrderID, OrderConfirmed)
				fmt.Println("[Order] status => CONFIRMED")
			case "RESERVE_FAILED":
				orderDB.UpdateOrderStatus(e.OrderID, OrderCanceled)
				fmt.Println("[Order] status => CANCELED")
			case "DUPLICATE":
				fmt.Println("[Order] duplicate consume, no status overwrite")
			}
		}
		close(done)
	}()

	<-done
	fmt.Println("final order:", orderDB.GetOrder("o-1001"))

	fmt.Println("\n===== compensation demo =====")
	order2 := Order{OrderID: "o-1002", SKU: "sku-1", Qty: 2, Status: OrderPending}
	event2 := Event{EventID: "evt-2", Type: "OrderCreated", OrderID: order2.OrderID, SKU: order2.SKU, Qty: order2.Qty}
	orderDB.CreateOrderWithOutbox(order2, event2)
	result := inventory.HandleOrderCreated(event2)
	if result == "RESERVE_FAILED" {
		orderDB.UpdateOrderStatus(order2.OrderID, OrderCanceled)
		fmt.Println("[Order] cancel directly because reserve failed")
	} else {
		fmt.Println("[Shipping] create shipment failed, trigger compensation")
		inventory.Release(order2.OrderID, order2.SKU, order2.Qty)
		orderDB.UpdateOrderStatus(order2.OrderID, OrderCanceled)
	}
	fmt.Println("final order2:", orderDB.GetOrder("o-1002"))
}

执行方式:

go run main.go

这个示例表达了几个非常关键的设计点:

  • 本地事务和消息投递解耦:先把业务数据和 outbox 一起提交
  • 消息可重试:发布失败可以继续扫表重发
  • 消费幂等:即使消息重复,也不会重复扣库存
  • 补偿机制:后置步骤失败时,可以把前面步骤撤销回来

这正是最终一致性方案能落地的核心。

五、分布式事务方案对比:2PC、Saga、TCC

当一个业务操作跨多个服务或多个资源管理器时,常见的分布式事务方案主要有 2PC、Saga、TCC。

5.1 2PC(Two-Phase Commit,两阶段提交)

2PC 把分布式事务分成两个阶段:

  1. Prepare 阶段:协调者询问所有参与者是否可以提交
  2. Commit / Rollback 阶段:如果全部参与者都准备成功,则统一提交;否则统一回滚

优点:

  • 理论上更接近强一致
  • 对事务语义比较直观

缺点:

  • 同步阻塞明显
  • 协调者压力大,存在单点风险
  • 参与者资源锁定时间长,吞吐低
  • 在微服务体系中通常不够灵活

适用场景:

  • 参与节点少
  • 强一致要求高
  • 资源管理器对 XA / 2PC 支持较好

5.2 Saga

Saga 的思路是:把一个长事务拆成多个本地事务,每个本地事务成功后进入下一步;如果后续某一步失败,则按相反顺序执行补偿动作。

例如:

  1. 创建订单
  2. 扣减库存
  3. 创建发货单
  4. 如果第 3 步失败,则补偿第 2 步恢复库存,再补偿第 1 步取消订单

优点:

  • 不需要长时间锁资源
  • 更适合微服务场景
  • 吞吐高,可扩展性好

缺点:

  • 业务补偿逻辑复杂
  • 补偿不一定是严格意义上的“回滚”
  • 系统会经历中间态

适用场景:

  • 电商、营销、履约等业务流程长、步骤多的场景
  • 可以接受最终一致
  • 各步骤具有明确补偿动作

5.3 TCC(Try-Confirm-Cancel)

TCC 是业务层面的两阶段模型:

  • Try:预留资源
  • Confirm:正式提交
  • Cancel:取消预留

比如下单时:

  • Try:冻结库存、冻结余额
  • Confirm:正式扣减库存、扣款
  • Cancel:释放冻结库存、解冻余额

优点:

  • 一致性比普通 Saga 更强
  • 资源状态更明确
  • 适合核心交易链路

缺点:

  • 侵入业务非常深
  • 每个服务都要实现 Try / Confirm / Cancel 三套接口
  • 开发和测试成本高

适用场景:

  • 支付、资金、库存冻结等核心交易场景
  • 对状态控制要求高
  • 团队有能力长期维护较复杂的事务模型

5.4 三种方案横向对比

方案 一致性强度 性能 侵入性 复杂度 适用业务
2PC 少节点、强一致
Saga 长流程、最终一致
TCC 较高 很高 支付、库存冻结、资金链路

5.5 如何选型

可以按下面的思路判断:

  • 必须强一致,且参与者少:优先评估 2PC / XA
  • 链路长、追求可用性和吞吐:优先 Saga
  • 关键资源要先冻结再确认:优先 TCC

但从现代微服务实践看,大多数业务系统最终都会走向这样一种组合:

  • 核心资金链路:TCC 或更严格的状态机控制
  • 一般业务链路:Saga + MQ + 补偿
  • 展示和同步类链路:最终一致性 + 异步修复

六、补偿机制与幂等设计

补偿和幂等,是分布式一致性设计里最容易被低估、但最决定成败的两个点。

6.1 为什么一定要有补偿机制

在分布式环境下,失败不是偶发,而是常态:

  • 网络超时
  • 服务重启
  • 消息重复投递
  • 下游部分成功
  • 重试造成副作用叠加

如果系统没有补偿机制,就会出现很多“卡住的中间态”:

  • 订单创建了,但库存没扣成功
  • 库存扣了,但发货单没创建
  • 优惠券发了,但订单最终失败

补偿机制的目标,不一定是“把世界恢复到绝对原点”,而是:把业务状态拉回到可接受、可解释、可继续处理的状态。

6.2 补偿设计的关键原则

原则一:补偿动作必须可重复执行

补偿本身也可能失败、超时、重试,所以补偿接口也必须幂等。

原则二:补偿要有明确边界

不是所有动作都能补偿。例如短信已发送、邮件已投递、外部三方已感知的动作,通常只能做“逆向补偿”,而不能真的回滚。

原则三:补偿依赖状态机,而不是 if-else 拼接

建议为关键业务对象设计明确状态流转,例如:

  • 订单:INIT -> PENDING -> CONFIRMED / CANCELED
  • 库存预留:NONE -> RESERVED -> CONFIRMED / RELEASED
  • 支付:INIT -> FROZEN -> PAID / UNFROZEN

状态机越明确,补偿越不容易写乱。

6.3 为什么消费端必须幂等

消息队列通常只能保证至少一次投递,这意味着:

  • 同一条消息可能被重复消费
  • 消费成功但 ACK 丢失时,消息可能再次投递
  • 消费者重启后可能重复处理历史消息

如果没有幂等控制,就可能发生:

  • 库存被扣两次
  • 积分加两次
  • 优惠券发多张
  • 状态被旧消息覆盖

6.4 幂等常见实现方式

方式 核心思路 适用场景
唯一业务键去重 例如订单号、支付单号只能成功一次 下单、支付、发券
消息 ID 去重表 消费前检查 message_id 是否已处理 MQ 消费端
状态机幂等 只允许合法状态迁移 订单、库存、支付状态流转
数据库唯一索引 用唯一约束防重复插入 创建型操作
乐观锁 / 版本号 防止旧写覆盖新写 并发更新场景

6.5 工程上非常实用的落地建议

  1. 所有关键消息都带唯一 ID
  2. 所有消费逻辑都先做去重判断
  3. 所有状态更新都通过状态机控制
  4. 补偿接口与正常接口一样要求幂等
  5. 重试要有上限,并接入死信队列或人工补单机制
  6. 日志里要能串起 traceId、orderId、messageId

七、实践中的设计建议

最后,把缓存、消息队列、事务一致性放在一起看,给出几条更贴近工程落地的建议。

7.1 不要把缓存当真相源

缓存的职责是加速访问,不是定义业务真相。业务真相应由数据库、账本、状态机或事件日志来承载。

7.2 优先把一致性问题分层处理

可以按这个顺序去设计:

  1. 单机内先保证本地事务正确
  2. 跨服务用消息实现最终一致
  3. 对失败链路补偿
  4. 对重复执行做幂等
  5. 对缓存层单独做失效和回填治理

这样会比“一上来就追求全链路强一致”更现实。

7.3 先设计失败路径,再设计成功路径

一个成熟系统的设计顺序应该是:

  • 如果数据库成功、发消息失败怎么办?
  • 如果消息投递成功、消费失败怎么办?
  • 如果消费成功、回 ACK 失败怎么办?
  • 如果补偿执行了一半又失败怎么办?
  • 如果缓存删失败怎么办?

把这些问题先想清楚,系统的可靠性通常就不会太差。

7.4 关键链路尽量状态化,而不是事件散落

越是复杂的业务,越应该让订单、库存、支付这些对象有明确状态,而不是依赖多个布尔值拼接判断。状态化设计会极大提升可观测性、排障能力和补偿正确性。

八、总结

缓存解决的是性能问题,但会带来一致性挑战;消息队列解决的是解耦和削峰问题,但不会自动帮你完成事务语义;分布式事务方案解决的是跨服务协同问题,但任何方案都需要面对失败、重试、重复、补偿这些工程现实。

因此,真正可落地的一致性设计,往往不是依赖某一种“银弹技术”,而是几种机制的组合:

  • Cache-Aside 处理大多数缓存场景
  • 空值缓存、互斥回源、TTL 抖动 处理缓存异常
  • MQ + Outbox 建立最终一致性链路
  • 幂等消费 + 状态机 控制重复执行风险
  • 补偿机制 兜底不可避免的局部失败
  • 根据业务要求在 2PC、Saga、TCC 之间做权衡

如果要用一句话概括本文的核心思想,那就是:一致性不是“某个组件”的能力,而是整个系统围绕失败场景做出的协同设计。


📝 版权声明:本文为原创技术博客,转载请注明出处。

如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!

上一篇

Golang工程化: 5.5 单体到微服务的架构演进

下一篇

Golang工程化:5.7 稳定性设计之超时、重试、熔断、幂等