返回首页

Golang高级:3.5 高级并发模式

当我们从“会用 goroutine 和 channel”继续往前走,就会发现并发编程真正困难的部分,已经不再是语法本身,而是在高吞吐、强竞争、跨实例协作、错误传播和状态一致性之间做工程化取舍

在业务系统里,下面这些问题比“如何启动一个 goroutine”更常见:

  • 热点请求瞬时打爆缓存和数据库,如何避免重复计算
  • 多个并发子任务中只要有一个失败,如何快速取消其余任务
  • 高频读写的共享状态,如何在锁开销、正确性和复杂度之间平衡
  • 单机内互斥已经不够时,多个服务实例之间如何争抢同一份资源

这一章聚焦五个高级并发主题:

  • 无锁编程与 sync/atomic
  • singleflight:请求合并
  • errgroup:错误组管理
  • 并发数据结构设计
  • 分布式锁的 Go 实现

为了方便建立整体认知,可以先看一个总览。

模式 解决的问题 适用范围 关键词
sync/atomic 避免简单共享状态上的锁竞争 单变量或少量原子状态 CAS、原子读写、无锁热点路径
singleflight 合并相同 key 的重复请求 单进程内重复计算抑制 缓存击穿、请求去重
errgroup 并发任务的统一收敛与错误传播 一组生命周期一致的子任务 取消传播、失败即终止
并发数据结构 高并发场景下的数据访问效率与一致性 单机内共享内存结构 分片、局部锁、快照
Redis 分布式锁 多实例之间的互斥访问 分布式部署场景 SET NX EX、唯一 token、Lua 解锁

一、无锁编程与 sync/atomic

1.1 为什么会需要原子操作

锁并不是坏东西。事实上,在绝大多数业务代码里,sync.Mutex 依然是最稳妥、最容易维护的选择。问题在于,有些热点状态非常简单,例如:

  • 请求计数器
  • 限流器中的令牌数量
  • 连接数、在线人数、命中次数
  • 某个状态位是否已初始化
  • 基于 CAS 的库存扣减或任务抢占

如果只是对一个整数做读、写、加减,或者对一个指针做安全发布,那么引入完整互斥锁往往会让路径变重。这个时候,sync/atomic 就很合适。

原子操作的核心价值有两个:

  1. 单个共享状态的更新不可分割
  2. 多个 goroutine 对该状态的观察顺序有统一保证

但也必须强调:原子操作不是“更高级的锁替代品”,它只是“更窄但更轻”的并发工具。 一旦你的业务不再是“单个变量的状态变更”,而是涉及多个字段之间的一致性约束,锁通常比原子操作更安全。

1.2 atomic 适合什么,不适合什么

适合场景 原因
简单计数器 Add/Load/Store 成本低,表达清晰
状态位切换 例如 0/1 状态、是否关闭、是否初始化
CAS 抢占 适合“读当前值 -> 验证 -> 尝试更新”的竞争逻辑
指针安全发布 例如热更新配置快照
不适合场景 原因
多字段一致性 多个原子变量组合后不自动形成事务
复杂临界区 CAS 循环会让逻辑变难读、难验证
涉及外部调用的临界区 原子操作无法保护 IO、RPC、文件等副作用
需要条件变量/等待队列 这不是 atomic 的职责

1.3 CAS 的本质:乐观并发控制

CAS(Compare-And-Swap)可以理解为:

  1. 先读取当前值。
  2. 如果当前值仍然等于我刚才读到的值,就把它替换成新值。
  3. 如果已经被别人改过,说明竞争失败,重新来一轮。

这是一种典型的乐观并发控制:默认大家不会频繁冲突,只有真的冲突了才重试。

下面给出一个完整可运行示例:使用 atomic.Int64 和 CAS 循环实现一个“无锁库存扣减器”。

package main

import (
	"fmt"
	"sync"
	"sync/atomic"
)

type Inventory struct {
	stock atomic.Int64
	sold  atomic.Int64
}

func NewInventory(initial int64) *Inventory {
	inv := &Inventory{}
	inv.stock.Store(initial)
	return inv
}

func (i *Inventory) TryPurchase(n int64) bool {
	for {
		current := i.stock.Load()
		if current < n {
			return false
		}

		if i.stock.CompareAndSwap(current, current-n) {
			i.sold.Add(n)
			return true
		}
	}
}

func main() {
	inventory := NewInventory(100)

	var wg sync.WaitGroup
	var successBuyers atomic.Int64

	buyers := 40
	unitsPerBuyer := int64(3)

	for id := 1; id <= buyers; id++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()

			if inventory.TryPurchase(unitsPerBuyer) {
				successBuyers.Add(1)
				fmt.Printf("buyer %02d purchase success\n", id)
				return
			}

			fmt.Printf("buyer %02d purchase failed: stock not enough\n", id)
		}(id)
	}

	wg.Wait()

	fmt.Println("------------------------------")
	fmt.Println("remaining stock:", inventory.stock.Load())
	fmt.Println("sold units:", inventory.sold.Load())
	fmt.Println("successful buyers:", successBuyers.Load())
}

这个例子值得注意的点有三个。

第一,TryPurchase 没有加锁,而是不断尝试完成“读当前库存 -> 验证是否足够 -> CAS 更新库存”。如果中途别的 goroutine 抢先改了库存,CompareAndSwap 会失败,于是进入下一轮。

第二,这种写法适合短小、纯内存、冲突可接受的逻辑。库存扣减本身非常短,所以 CAS 自旋成本通常低于把整个路径放进一把互斥锁里。

第三,sold 虽然也是原子变量,但它与 stock 并不构成真正的事务。也就是说,如果你的业务要求“库存、销售记录、订单状态”三者必须同时一致,那么单靠原子操作远远不够,必须引入更强的一致性控制手段。

1.4 使用 atomic 的工程注意事项

注意点 说明
不要把多个原子变量误当作事务 a.Store()b.Store() 之间仍然可能被其他 goroutine 观察到中间态
CAS 循环内不要做重操作 包括日志风暴、复杂分配、IO 调用,否则重试成本很高
优先使用新类型封装 atomic.Int64atomic.Boolatomic.Pointer[T] 比旧式函数接口更直观
先追求正确,再追求无锁 如果锁版本已经足够快,没必要为了“高级”而强行原子化

可以把 atomic 理解为:用非常低的同步成本,换取非常有限但高频的并发安全能力。它擅长处理热点路径上的“小状态”,不擅长处理复杂业务事务。

二、singleflight:请求合并

2.1 它解决的不是缓存命中,而是缓存击穿时的重复计算

很多人第一次看到 singleflight 时,会误以为它是缓存组件。其实不是。它不负责缓存数据,而是负责当多个 goroutine 在同一时刻请求同一个 key 时,只让其中一个真正执行,其余请求复用同一份结果

这类问题最经典的场景就是缓存击穿

假设商品详情本来都在本地缓存里。某一时刻,热点商品 sku:1001 的缓存过期了。此时 1000 个请求同时到来:

  1. 所有请求先查缓存,发现 miss。
  2. 如果没有额外控制,这 1000 个请求会同时打到数据库。
  3. 数据库瞬间承压,延迟上升,甚至引发更大范围雪崩。

这时,singleflight 的价值就体现出来了:相同 key 的 miss 只放行一次真正加载,其余请求等待这次加载完成并共享结果。

2.2 singleflight 的边界

在使用它之前,需要先明确边界。

结论 说明
它不能替代缓存 它只是“合并同一时刻的重复请求”
它通常作用于单进程内 多实例之间不能靠它做全局去重
key 设计非常关键 key 太粗会误合并,key 太细会失去效果
适合高代价加载路径 例如 DB 查询、远程计算、模型推理、配置拉取

2.3 缓存击穿场景下的完整示例

下面的程序演示一个简化版商品详情服务:

  • 先查本地缓存
  • 缓存 miss 时,通过 singleflight.Group 合并相同商品 ID 的加载请求
  • 模拟数据库加载耗时 200ms
  • 最终观察真实加载次数只发生一次
package main

import (
	"context"
	"fmt"
	"sync"
	"sync/atomic"
	"time"

	"golang.org/x/sync/singleflight"
)

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

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

func (c *Cache) Get(key string) (string, bool) {
	c.mu.RLock()
	defer c.mu.RUnlock()
	v, ok := c.data[key]
	return v, ok
}

func (c *Cache) Set(key, value string) {
	c.mu.Lock()
	defer c.mu.Unlock()
	c.data[key] = value
}

type ProductService struct {
	cache     *Cache
	group     singleflight.Group
	loadCount atomic.Int64
}

func NewProductService() *ProductService {
	return &ProductService{cache: NewCache()}
}

func (s *ProductService) loadFromDB(ctx context.Context, productID string) (string, error) {
	s.loadCount.Add(1)
	fmt.Println("load from DB for", productID)

	select {
	case <-time.After(200 * time.Millisecond):
		return "detail-of-" + productID, nil
	case <-ctx.Done():
		return "", ctx.Err()
	}
}

func (s *ProductService) GetProduct(ctx context.Context, productID string) (string, error) {
	if value, ok := s.cache.Get(productID); ok {
		return value, nil
	}

	result, err, shared := s.group.Do(productID, func() (any, error) {
		if value, ok := s.cache.Get(productID); ok {
			return value, nil
		}

		value, err := s.loadFromDB(ctx, productID)
		if err != nil {
			return "", err
		}

		s.cache.Set(productID, value)
		return value, nil
	})
	if err != nil {
		return "", err
	}

	fmt.Printf("productID=%s shared=%v\n", productID, shared)
	return result.(string), nil
}

func main() {
	service := NewProductService()
	ctx := context.Background()

	var wg sync.WaitGroup
	start := make(chan struct{})

	for i := 0; i < 10; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			<-start

			value, err := service.GetProduct(ctx, "sku:1001")
			if err != nil {
				fmt.Printf("worker %d error: %v\n", id, err)
				return
			}

			fmt.Printf("worker %d got %s\n", id, value)
		}(i)
	}

	close(start)
	wg.Wait()

	fmt.Println("------------------------------")
	fmt.Println("real DB load count:", service.loadCount.Load())

	value, err := service.GetProduct(ctx, "sku:1001")
	if err != nil {
		fmt.Println("second round error:", err)
		return
	}
	fmt.Println("second round cache hit:", value)
}

这段代码的关键不只是 group.Do(productID, fn) 这一行,而是整个加载路径的分层:

  1. 先查缓存,命中则直接返回。
  2. 缓存 miss 后再进入 singleflight
  3. fn再次检查缓存,防止等待期间别的 goroutine 已经写入结果。
  4. 只有一个 goroutine 会真正执行 loadFromDB
  5. 结果写入缓存后,其他等待者直接复用。

2.4 singleflight 的实践要点

要点 说明
key 要能准确表达“相同请求” 例如商品详情可用商品 ID,搜索结果则可能需要拼接参数摘要
singleflight 只合并飞行中的请求 一旦第一个请求完成,下一波请求会重新触发
对错误结果要谨慎 如果加载失败,所有等待者会一起拿到同一个错误
和缓存搭配使用效果最好 单独使用只能去重,不能减少后续重复访问

简单概括就是:缓存负责“记住结果”,singleflight 负责“控制 miss 时不要同时算很多次”。

三、errgroup:错误组管理

3.1 为什么 WaitGroup 不够

sync.WaitGroup 能解决“等一组 goroutine 全部结束”,但它并不解决下面几个关键问题:

  • 子任务返回错误后,错误怎么统一收集
  • 某个子任务已经失败,其余任务要不要继续执行
  • 父任务超时或取消时,所有子任务如何及时停止

如果你自己在 WaitGroup 之外再叠加错误 channel、取消逻辑、状态收敛,很快就会把代码写得很散。errgroup 的价值,就在于把这些需求组合成一个更适合业务开发的并发抽象。

3.2 errgroup 的核心模型

errgroup 可以理解为“带错误传播和取消能力的 goroutine group”。最常见的使用方式是:

  1. errgroup.WithContext 创建一个 group 和派生上下文。
  2. g.Go 启动多个并发子任务。
  3. 任意一个任务返回错误,派生 ctx 会被取消。
  4. g.Wait() 返回第一个非空错误。

这对于“并发拉取多个下游数据,只要有一个失败就整体失败”的聚合接口非常适合。

3.3 完整示例:并发聚合用户视图

下面示例模拟一个聚合接口,同时并发查询:

  • 用户画像
  • 订单摘要
  • 账户余额

其中订单服务故意返回错误,看看其余任务如何被取消。

package main

import (
	"context"
	"errors"
	"fmt"
	"time"

	"golang.org/x/sync/errgroup"
)

func fetchProfile(ctx context.Context) (string, error) {
	select {
	case <-time.After(200 * time.Millisecond):
		fmt.Println("profile service done")
		return "profile: VIP user", nil
	case <-ctx.Done():
		fmt.Println("profile service canceled")
		return "", ctx.Err()
	}
}

func fetchOrders(ctx context.Context) (string, error) {
	select {
	case <-time.After(300 * time.Millisecond):
		fmt.Println("orders service failed")
		return "", errors.New("orders service timeout")
	case <-ctx.Done():
		fmt.Println("orders service canceled")
		return "", ctx.Err()
	}
}

func fetchBalance(ctx context.Context) (string, error) {
	select {
	case <-time.After(2 * time.Second):
		fmt.Println("balance service done")
		return "balance: 1024", nil
	case <-ctx.Done():
		fmt.Println("balance service canceled")
		return "", ctx.Err()
	}
}

func main() {
	g, ctx := errgroup.WithContext(context.Background())

	var profile string
	var orders string
	var balance string

	g.Go(func() error {
		result, err := fetchProfile(ctx)
		if err != nil {
			return err
		}
		profile = result
		return nil
	})

	g.Go(func() error {
		result, err := fetchOrders(ctx)
		if err != nil {
			return err
		}
		orders = result
		return nil
	})

	g.Go(func() error {
		result, err := fetchBalance(ctx)
		if err != nil {
			return err
		}
		balance = result
		return nil
	})

	if err := g.Wait(); err != nil {
		fmt.Println("aggregate request failed:", err)
		return
	}

	fmt.Println("aggregate request success")
	fmt.Println(profile)
	fmt.Println(orders)
	fmt.Println(balance)
}

运行这段程序时,你会观察到以下行为:

  • profile 很快成功返回。
  • orders 在 300ms 后失败。
  • balance 原本需要 2 秒,但因为 group 上下文被取消,所以提前结束。
  • 最终 g.Wait() 返回订单服务错误,整个聚合请求失败收敛。

这正是 errgroup 最常见的业务语义:一组任务共同服务于同一个目标,只要关键路径中有一个失败,其余任务就没有继续执行的必要。

3.4 errgroup 的使用注意事项

注意点 说明
子任务必须尊重 ctx.Done() 否则 group 取消了,任务还是会继续跑
共享结果写入要考虑并发安全 上例中每个变量只被一个 goroutine 写,因此安全
Wait() 返回第一个错误 它不是完整错误聚合器,如果你要收集所有错误,需要额外设计
非关键失败不要滥用“失败即取消” 有些场景更适合局部降级,而不是整体失败

如果说 WaitGroup 只是“计数器”,那么 errgroup 就是更贴近业务的“并发任务生命周期管理器”。

四、并发数据结构设计

4.1 真正困难的不是“加锁”,而是“设计访问模型”

谈并发数据结构时,很多人第一反应是“给 map 加一把锁”。这只能解决最粗粒度的线程安全问题,却不一定能解决吞吐、延迟、扩展性和接口语义问题。

一个并发数据结构设计得好不好,通常取决于下面几件事:

设计维度 关键问题
数据归属 哪些 goroutine 会访问它,访问模式是否集中
读写比例 是读多写少,还是高频更新
一致性要求 是否需要强一致快照,还是最终可接受
锁粒度 是全局锁、分段锁,还是无锁热点字段
API 语义 是否把复合操作封装成原子接口

换句话说,并发数据结构首先是“数据与访问模式的设计题”,其次才是“用什么同步原语”的实现题。

4.2 典型思路:分片(Sharding)

如果你有一个高并发 map,最简单的做法是:

  • 整个 map 外面套一把 RWMutex
  • 读时加读锁,写时加写锁

这在中低并发场景往往完全够用。但如果 key 很多、写入也频繁,单把全局锁容易形成热点。一个常见优化思路就是分片

  1. 按 hash 把 key 打散到多个 shard。
  2. 每个 shard 自己维护一把锁和一个局部 map。
  3. 不同 shard 的操作可以并行进行。

这本质上是用“结构复杂度”换“更低的锁竞争”。

4.3 完整示例:分片计数器 Map

下面实现一个可运行的分片计数器,用于统计接口访问次数。它支持:

  • Add(key, delta):累加计数
  • Get(key):读取单个 key
  • Snapshot():导出全量快照
package main

import (
	"fmt"
	"hash/fnv"
	"sync"
)

const shardCount = 16

type counterShard struct {
	mu   sync.RWMutex
	data map[string]int
}

type ShardedCounterMap struct {
	shards [shardCount]counterShard
}

func NewShardedCounterMap() *ShardedCounterMap {
	m := &ShardedCounterMap{}
	for i := range m.shards {
		m.shards[i].data = make(map[string]int)
	}
	return m
}

func (m *ShardedCounterMap) shardFor(key string) *counterShard {
	h := fnv.New32a()
	_, _ = h.Write([]byte(key))
	idx := h.Sum32() % shardCount
	return &m.shards[idx]
}

func (m *ShardedCounterMap) Add(key string, delta int) {
	shard := m.shardFor(key)
	shard.mu.Lock()
	defer shard.mu.Unlock()
	shard.data[key] += delta
}

func (m *ShardedCounterMap) Get(key string) int {
	shard := m.shardFor(key)
	shard.mu.RLock()
	defer shard.mu.RUnlock()
	return shard.data[key]
}

func (m *ShardedCounterMap) Snapshot() map[string]int {
	result := make(map[string]int)
	for i := range m.shards {
		shard := &m.shards[i]
		shard.mu.RLock()
		for k, v := range shard.data {
			result[k] = v
		}
		shard.mu.RUnlock()
	}
	return result
}

func main() {
	counter := NewShardedCounterMap()

	endpoints := []string{
		"/api/user/profile",
		"/api/order/list",
		"/api/feed/recommend",
		"/api/cart/add",
	}

	var wg sync.WaitGroup

	for i := 0; i < 100; i++ {
		for _, endpoint := range endpoints {
			wg.Add(1)
			go func(endpoint string) {
				defer wg.Done()
				counter.Add(endpoint, 1)
			}(endpoint)
		}
	}

	wg.Wait()

	fmt.Println("/api/user/profile =>", counter.Get("/api/user/profile"))
	fmt.Println("/api/order/list =>", counter.Get("/api/order/list"))
	fmt.Println("/api/feed/recommend =>", counter.Get("/api/feed/recommend"))
	fmt.Println("/api/cart/add =>", counter.Get("/api/cart/add"))

	fmt.Println("------------------------------")
	fmt.Println("snapshot:")
	for k, v := range counter.Snapshot() {
		fmt.Printf("%s => %d\n", k, v)
	}
}

4.4 这个设计为什么比“全局一把锁”更进一步

这个例子并不复杂,但它体现了并发数据结构设计里几个很关键的原则。

原则一:把竞争拆散

不同 endpoint 会被 hash 到不同 shard。这样即使很多 goroutine 同时写入,也不一定都挤在同一把锁上。

原则二:把复合操作收进结构内部

调用方不需要自己“先查 map,再加锁,再修改”。它只调用 AddGetSnapshot 这些明确语义的方法。这一点非常重要:线程安全不应该依赖调用方自觉,而应该尽量由数据结构自身保证。

原则三:快照要明确语义边界

Snapshot() 是逐 shard 拿读锁并复制数据的,因此它拿到的是一个“近似一致”的导出结果,而不是全局时刻完全冻结的事务快照。对于监控统计、指标导出,这通常足够;但对金融账务类业务则不够。

4.5 并发数据结构的常见取舍

方案 优点 缺点 适合场景
全局 Mutex/RWMutex 实现简单,容易验证正确性 高并发下可能形成热点 中低并发、逻辑清晰优先
分片锁结构 降低锁竞争,扩展性更好 实现更复杂,快照语义更弱 高频读写、key 分布较散
sync.Map 使用方便,免去模板代码 类型不直观,复合操作表达弱 读多写少、键集合动态变化
原子字段 + 不可变快照 读路径非常快 写路径实现更难 配置、路由表、只读快照发布

真正好的并发数据结构,从来不是“锁越少越好”,而是:在正确性前提下,尽量让锁竞争与业务访问模式匹配。

五、分布式锁的 Go 实现

5.1 为什么单机锁不够

sync.MutexRWMutexatomic 都只能在单个进程内部保证并发安全。如果服务已经多实例部署,那么下面这种问题就会出现:

  • 定时任务只能被一个实例执行
  • 同一批数据只能被一个 worker 抢到
  • 同一个用户的幂等处理不能跨实例重复执行
  • 某个关键资源更新必须在集群范围内串行化

这个时候,单机锁已经失效,因为实例之间根本不共享内存。于是我们需要一种基于外部共享存储的互斥机制,常见选择就是 Redis 分布式锁。

5.2 Redis 锁的基本思路

一个相对可靠的最小实现,至少要满足下面两点:

  1. 加锁使用原子命令:通常是 SET key value NX EX ttl
  2. 解锁必须校验持有者身份:只有 value 等于自己 token 时,才能删除 key。

第二点尤其重要。很多错误实现会直接 DEL key,这会导致:

  • 实例 A 的锁过期了
  • 实例 B 获得了同一把锁
  • 这时实例 A 才执行解锁,如果直接 DEL key,就会把 B 的锁误删掉

因此,解锁必须通过 Lua 脚本做“比较 value 再删除”的原子操作。

5.3 完整示例:基于 Redis 的 Go 分布式锁

下面给出一个完整示例,依赖 github.com/redis/go-redis/v9。运行前请确保本地 Redis 可用,例如监听在 127.0.0.1:6379

package main

import (
	"context"
	"crypto/rand"
	"encoding/hex"
	"errors"
	"fmt"
	"time"

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

var unlockScript = redis.NewScript(`
if redis.call("get", KEYS[1]) == ARGV[1] then
    return redis.call("del", KEYS[1])
end
return 0
`)

type RedisLock struct {
	client     *redis.Client
	key        string
	token      string
	expiration time.Duration
}

func NewRedisLock(client *redis.Client, key string, expiration time.Duration) (*RedisLock, error) {
	token, err := randomToken()
	if err != nil {
		return nil, err
	}
	return &RedisLock{
		client:     client,
		key:        key,
		token:      token,
		expiration: expiration,
	}, nil
}

func randomToken() (string, error) {
	buf := make([]byte, 16)
	if _, err := rand.Read(buf); err != nil {
		return "", err
	}
	return hex.EncodeToString(buf), nil
}

func (l *RedisLock) Acquire(ctx context.Context, retryInterval, maxWait time.Duration) error {
	deadline := time.Now().Add(maxWait)

	for {
		ok, err := l.client.SetNX(ctx, l.key, l.token, l.expiration).Result()
		if err != nil {
			return err
		}
		if ok {
			return nil
		}

		if time.Now().After(deadline) {
			return errors.New("acquire lock timeout")
		}

		select {
		case <-time.After(retryInterval):
		case <-ctx.Done():
			return ctx.Err()
		}
	}
}

func (l *RedisLock) Release(ctx context.Context) error {
	result, err := unlockScript.Run(ctx, l.client, []string{l.key}, l.token).Int()
	if err != nil {
		return err
	}
	if result == 0 {
		return errors.New("lock already expired or not owned by current client")
	}
	return nil
}

func worker(ctx context.Context, client *redis.Client, name string, hold time.Duration) {
	lock, err := NewRedisLock(client, "lock:order:close:job", 3*time.Second)
	if err != nil {
		fmt.Printf("%s create lock failed: %v\n", name, err)
		return
	}

	fmt.Printf("%s try acquire lock\n", name)
	if err := lock.Acquire(ctx, 200*time.Millisecond, 5*time.Second); err != nil {
		fmt.Printf("%s acquire failed: %v\n", name, err)
		return
	}

	fmt.Printf("%s acquired lock, start job\n", name)
	defer func() {
		if err := lock.Release(ctx); err != nil {
			fmt.Printf("%s release failed: %v\n", name, err)
			return
		}
		fmt.Printf("%s released lock\n", name)
	}()

	time.Sleep(hold)
	fmt.Printf("%s finished job\n", name)
}

func main() {
	ctx := context.Background()
	client := redis.NewClient(&redis.Options{
		Addr: "127.0.0.1:6379",
	})
	defer client.Close()

	if err := client.Ping(ctx).Err(); err != nil {
		fmt.Println("ping redis failed:", err)
		return
	}

	go worker(ctx, client, "instance-A", 2*time.Second)
	time.Sleep(100 * time.Millisecond)
	go worker(ctx, client, "instance-B", 1*time.Second)

	time.Sleep(6 * time.Second)
}

这个示例展示了最关键的两个动作:

  • 加锁SetNX(key, token, expiration),只有 key 不存在时才能成功。
  • 解锁:Lua 脚本先比较 token,再决定是否删除 key。

instance-A 获得锁后,instance-B 会不断重试,直到 A 执行完成并释放锁,然后 B 才有机会进入临界区。

5.4 这个实现仍然有哪些边界

虽然上面的实现已经比“直接 SETNX + DEL”可靠得多,但在生产环境里仍然要继续考虑以下问题。

问题 说明
锁超时 如果业务执行时间超过 expiration,锁可能提前过期,被别的实例拿走
锁续租 长任务往往需要 watchdog 或定时续期机制
时钟与暂停 GC stop-the-world、容器调度暂停都可能影响持锁时长判断
业务幂等 分布式锁只能降低并发冲突,不应替代最终业务幂等设计
高可用争议 极端故障下是否满足严格一致性,要看你的容错要求和部署架构

所以更准确地说,Redis 分布式锁通常适合做:

  • 定时任务单实例执行
  • 短事务型互斥控制
  • 抢占式 worker 协调
  • 对偶发重复执行有一定容忍度,但希望显著降低冲突概率的场景

如果业务需要的是强一致的分布式协调,往往还需要结合数据库唯一约束、幂等表、版本号、租约机制,甚至使用专门的协调系统,而不是把所有一致性问题都压到一把 Redis 锁上。

六、如何把这些模式用在真实系统里

到这里,你会发现这些高级并发模式并不是彼此割裂的,它们经常会组合出现。

例如一个高并发接口的常见设计可能是:

  1. 入口请求先查本地缓存。
  2. 缓存 miss 后用 singleflight 合并同 key 加载。
  3. 加载逻辑内部再用 errgroup 并发请求多个下游。
  4. 热点统计和运行状态通过 atomic 维护。
  5. 如果需要跨实例串行执行某个后台修复任务,再引入 Redis 分布式锁。

可以用下面这张表做一个选择参考。

问题场景 优先考虑的方案
高频计数、状态位、轻量 CAS 更新 sync/atomic
缓存击穿导致同 key 重复回源 singleflight
一组并发子任务需要统一取消和错误返回 errgroup
本地共享结构读写竞争明显 分片锁 / 自定义并发数据结构
多实例之间要互斥执行任务 Redis 分布式锁

高级并发的核心,不是记住多少 API,而是建立一个稳定判断标准:

  • 这是单机问题,还是分布式问题?
  • 这是状态同步问题,还是任务生命周期问题?
  • 我需要的是强一致,还是吞吐优先?
  • 我的同步开销,是不是已经超过业务本身?

想清楚这些问题,你就不会陷入“见并发就上 channel,见竞争就上锁,见高阶就盲目无锁”的误区。

七、小结

本章的五个主题,本质上分别回答了五类不同问题:

  • sync/atomic:如何以更低成本维护简单共享状态
  • singleflight:如何在缓存击穿时合并同类请求
  • errgroup:如何让一组并发任务以“业务整体”视角收敛
  • 并发数据结构设计:如何让共享内存结构与访问模式匹配
  • Redis 分布式锁:如何把互斥语义从单进程扩展到多实例

如果说基础并发关注的是“怎么并发起来”,那么高级并发更关注的是:怎么在竞争、失败、扩展和一致性之间做出有代价意识的设计。

真正成熟的 Go 并发代码,往往不是用了多少技巧,而是知道什么时候该简单,什么时候该重武器上场,什么时候又必须承认“这已经不是单机并发问题,而是分布式一致性问题”。


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

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

上一篇

Golang高级:3.4 Unsafe 与底层编程

下一篇

Golang高级:3.6 编译器与链接器