当我们从“会用 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 就很合适。
原子操作的核心价值有两个:
- 单个共享状态的更新不可分割。
- 多个 goroutine 对该状态的观察顺序有统一保证。
但也必须强调:原子操作不是“更高级的锁替代品”,它只是“更窄但更轻”的并发工具。 一旦你的业务不再是“单个变量的状态变更”,而是涉及多个字段之间的一致性约束,锁通常比原子操作更安全。
1.2 atomic 适合什么,不适合什么
| 适合场景 | 原因 |
|---|---|
| 简单计数器 | Add/Load/Store 成本低,表达清晰 |
| 状态位切换 | 例如 0/1 状态、是否关闭、是否初始化 |
| CAS 抢占 | 适合“读当前值 -> 验证 -> 尝试更新”的竞争逻辑 |
| 指针安全发布 | 例如热更新配置快照 |
| 不适合场景 | 原因 |
|---|---|
| 多字段一致性 | 多个原子变量组合后不自动形成事务 |
| 复杂临界区 | CAS 循环会让逻辑变难读、难验证 |
| 涉及外部调用的临界区 | 原子操作无法保护 IO、RPC、文件等副作用 |
| 需要条件变量/等待队列 | 这不是 atomic 的职责 |
1.3 CAS 的本质:乐观并发控制
CAS(Compare-And-Swap)可以理解为:
- 先读取当前值。
- 如果当前值仍然等于我刚才读到的值,就把它替换成新值。
- 如果已经被别人改过,说明竞争失败,重新来一轮。
这是一种典型的乐观并发控制:默认大家不会频繁冲突,只有真的冲突了才重试。
下面给出一个完整可运行示例:使用 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.Int64、atomic.Bool、atomic.Pointer[T] 比旧式函数接口更直观 |
| 先追求正确,再追求无锁 | 如果锁版本已经足够快,没必要为了“高级”而强行原子化 |
可以把 atomic 理解为:用非常低的同步成本,换取非常有限但高频的并发安全能力。它擅长处理热点路径上的“小状态”,不擅长处理复杂业务事务。
二、singleflight:请求合并
2.1 它解决的不是缓存命中,而是缓存击穿时的重复计算
很多人第一次看到 singleflight 时,会误以为它是缓存组件。其实不是。它不负责缓存数据,而是负责当多个 goroutine 在同一时刻请求同一个 key 时,只让其中一个真正执行,其余请求复用同一份结果。
这类问题最经典的场景就是缓存击穿。
假设商品详情本来都在本地缓存里。某一时刻,热点商品 sku:1001 的缓存过期了。此时 1000 个请求同时到来:
- 所有请求先查缓存,发现 miss。
- 如果没有额外控制,这 1000 个请求会同时打到数据库。
- 数据库瞬间承压,延迟上升,甚至引发更大范围雪崩。
这时,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) 这一行,而是整个加载路径的分层:
- 先查缓存,命中则直接返回。
- 缓存 miss 后再进入
singleflight。 - 在
fn里再次检查缓存,防止等待期间别的 goroutine 已经写入结果。 - 只有一个 goroutine 会真正执行
loadFromDB。 - 结果写入缓存后,其他等待者直接复用。
2.4 singleflight 的实践要点
| 要点 | 说明 |
|---|---|
| key 要能准确表达“相同请求” | 例如商品详情可用商品 ID,搜索结果则可能需要拼接参数摘要 |
singleflight 只合并飞行中的请求 |
一旦第一个请求完成,下一波请求会重新触发 |
| 对错误结果要谨慎 | 如果加载失败,所有等待者会一起拿到同一个错误 |
| 和缓存搭配使用效果最好 | 单独使用只能去重,不能减少后续重复访问 |
简单概括就是:缓存负责“记住结果”,singleflight 负责“控制 miss 时不要同时算很多次”。
三、errgroup:错误组管理
3.1 为什么 WaitGroup 不够
sync.WaitGroup 能解决“等一组 goroutine 全部结束”,但它并不解决下面几个关键问题:
- 子任务返回错误后,错误怎么统一收集
- 某个子任务已经失败,其余任务要不要继续执行
- 父任务超时或取消时,所有子任务如何及时停止
如果你自己在 WaitGroup 之外再叠加错误 channel、取消逻辑、状态收敛,很快就会把代码写得很散。errgroup 的价值,就在于把这些需求组合成一个更适合业务开发的并发抽象。
3.2 errgroup 的核心模型
errgroup 可以理解为“带错误传播和取消能力的 goroutine group”。最常见的使用方式是:
- 用
errgroup.WithContext创建一个 group 和派生上下文。 - 用
g.Go启动多个并发子任务。 - 任意一个任务返回错误,派生
ctx会被取消。 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 很多、写入也频繁,单把全局锁容易形成热点。一个常见优化思路就是分片:
- 按 hash 把 key 打散到多个 shard。
- 每个 shard 自己维护一把锁和一个局部 map。
- 不同 shard 的操作可以并行进行。
这本质上是用“结构复杂度”换“更低的锁竞争”。
4.3 完整示例:分片计数器 Map
下面实现一个可运行的分片计数器,用于统计接口访问次数。它支持:
Add(key, delta):累加计数Get(key):读取单个 keySnapshot():导出全量快照
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,再加锁,再修改”。它只调用 Add、Get、Snapshot 这些明确语义的方法。这一点非常重要:线程安全不应该依赖调用方自觉,而应该尽量由数据结构自身保证。
原则三:快照要明确语义边界
Snapshot() 是逐 shard 拿读锁并复制数据的,因此它拿到的是一个“近似一致”的导出结果,而不是全局时刻完全冻结的事务快照。对于监控统计、指标导出,这通常足够;但对金融账务类业务则不够。
4.5 并发数据结构的常见取舍
| 方案 | 优点 | 缺点 | 适合场景 |
|---|---|---|---|
全局 Mutex/RWMutex |
实现简单,容易验证正确性 | 高并发下可能形成热点 | 中低并发、逻辑清晰优先 |
| 分片锁结构 | 降低锁竞争,扩展性更好 | 实现更复杂,快照语义更弱 | 高频读写、key 分布较散 |
sync.Map |
使用方便,免去模板代码 | 类型不直观,复合操作表达弱 | 读多写少、键集合动态变化 |
| 原子字段 + 不可变快照 | 读路径非常快 | 写路径实现更难 | 配置、路由表、只读快照发布 |
真正好的并发数据结构,从来不是“锁越少越好”,而是:在正确性前提下,尽量让锁竞争与业务访问模式匹配。
五、分布式锁的 Go 实现
5.1 为什么单机锁不够
sync.Mutex、RWMutex、atomic 都只能在单个进程内部保证并发安全。如果服务已经多实例部署,那么下面这种问题就会出现:
- 定时任务只能被一个实例执行
- 同一批数据只能被一个 worker 抢到
- 同一个用户的幂等处理不能跨实例重复执行
- 某个关键资源更新必须在集群范围内串行化
这个时候,单机锁已经失效,因为实例之间根本不共享内存。于是我们需要一种基于外部共享存储的互斥机制,常见选择就是 Redis 分布式锁。
5.2 Redis 锁的基本思路
一个相对可靠的最小实现,至少要满足下面两点:
- 加锁使用原子命令:通常是
SET key value NX EX ttl。 - 解锁必须校验持有者身份:只有 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 锁上。
六、如何把这些模式用在真实系统里
到这里,你会发现这些高级并发模式并不是彼此割裂的,它们经常会组合出现。
例如一个高并发接口的常见设计可能是:
- 入口请求先查本地缓存。
- 缓存 miss 后用
singleflight合并同 key 加载。 - 加载逻辑内部再用
errgroup并发请求多个下游。 - 热点统计和运行状态通过
atomic维护。 - 如果需要跨实例串行执行某个后台修复任务,再引入 Redis 分布式锁。
可以用下面这张表做一个选择参考。
| 问题场景 | 优先考虑的方案 |
|---|---|
| 高频计数、状态位、轻量 CAS 更新 | sync/atomic |
| 缓存击穿导致同 key 重复回源 | singleflight |
| 一组并发子任务需要统一取消和错误返回 | errgroup |
| 本地共享结构读写竞争明显 | 分片锁 / 自定义并发数据结构 |
| 多实例之间要互斥执行任务 | Redis 分布式锁 |
高级并发的核心,不是记住多少 API,而是建立一个稳定判断标准:
- 这是单机问题,还是分布式问题?
- 这是状态同步问题,还是任务生命周期问题?
- 我需要的是强一致,还是吞吐优先?
- 我的同步开销,是不是已经超过业务本身?
想清楚这些问题,你就不会陷入“见并发就上 channel,见竞争就上锁,见高阶就盲目无锁”的误区。
七、小结
本章的五个主题,本质上分别回答了五类不同问题:
sync/atomic:如何以更低成本维护简单共享状态singleflight:如何在缓存击穿时合并同类请求errgroup:如何让一组并发任务以“业务整体”视角收敛- 并发数据结构设计:如何让共享内存结构与访问模式匹配
- Redis 分布式锁:如何把互斥语义从单进程扩展到多实例
如果说基础并发关注的是“怎么并发起来”,那么高级并发更关注的是:怎么在竞争、失败、扩展和一致性之间做出有代价意识的设计。
真正成熟的 Go 并发代码,往往不是用了多少技巧,而是知道什么时候该简单,什么时候该重武器上场,什么时候又必须承认“这已经不是单机并发问题,而是分布式一致性问题”。
📝 版权声明:本文为原创技术博客,转载请注明出处。
如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!