网站后端通知系统:无需互斥锁的Go语言PubSub实现方案问询
解决Go PubSub系统中高并发订阅/退订的锁竞争问题
首先,你的问题非常典型——当PubSub系统面临大量频道和频繁的订阅者变更时,粗粒度的互斥锁确实会成为性能瓶颈。下面我会分享几个在Go生态中常用的优化方案,既能减少锁竞争,又能保证通知的可靠性:
方案1:用atomic.Value实现无锁的订阅者列表更新
核心思路是避免在消息分发期间持有锁,通过原子操作替换整个订阅者切片,让读操作(分发消息)完全无锁,写操作(增删订阅者)只在替换切片的瞬间有极短的开销。
修改后的结构体和核心逻辑:
import ( "sync/atomic" "time" ) type Channel struct { ingress <-chan int // 接收新文章ID的频道 subs atomic.Value // 存储[]chan<- int,原子更新 } func NewChannel(ingress <-chan int) *Channel { c := &Channel{ingress: ingress} c.subs.Store([]chan<- int{}) // 初始化空切片 go c.run() return c } func (c *Channel) run() { for msg := range c.ingress { // 读取当前订阅者列表(无锁,读取的是快照) subs := c.subs.Load().([]chan<- int) // 遍历快照发送消息,此时不持有任何锁 for _, sub := range subs { select { case sub <- msg: // 消息发送成功 case <-time.After(50 * time.Millisecond): // 订阅者无响应,异步移除避免阻塞分发 go c.Unsubscribe(sub) } } } } func (c *Channel) Subscribe(sub chan<- int) { // 加载当前切片,复制后添加新订阅者 current := c.subs.Load().([]chan<- int) updated := append(current, sub) c.subs.Store(updated) } func (c *Channel) Unsubscribe(sub chan<- int) { current := c.subs.Load().([]chan<- int) updated := make([]chan<- int, 0, len(current)) for _, s := range current { if s != sub { updated = append(updated, s) } } c.subs.Store(updated) }
优缺点:
- ✅ 完全消除了分发消息时的锁竞争,读操作零开销
- ✅ 实现简单,不需要复杂的锁管理
- ❌ 当订阅者数量极大时(比如10w+),增删操作需要复制整个切片,会带来一定的内存和CPU开销
方案2:分段锁(Sharded Locks)降低锁粒度
如果你的系统订阅者数量特别庞大,atomic.Value的切片复制开销不可接受,可以尝试把订阅者列表拆分成多个分片,每个分片用独立的互斥锁保护。这样增删订阅者时只需要锁定一个分片,分发消息时并行遍历多个分片,锁竞争的概率会大幅降低。
修改后的结构体和核心逻辑:
import ( "sync" "time" "unsafe" ) const shardCount = 16 // 分片数量,可根据并发量调整 type Channel struct { ingress <-chan int shards [shardCount]struct { subs []chan<- int mx sync.Mutex } } func NewChannel(ingress <-chan int) *Channel { c := &Channel{ingress: ingress} go c.run() return c } // 根据订阅者channel的地址哈希选择分片 func (c *Channel) getShard(sub chan<- int) *struct { subs []chan<- int mx sync.Mutex } { hash := uintptr(unsafe.Pointer(&sub)) % shardCount return &c.shards[hash] } func (c *Channel) run() { for msg := range c.ingress { // 遍历所有分片,并行发送消息(也可串行,根据性能需求调整) var wg sync.WaitGroup for i := range c.shards { wg.Add(1) go func(i int) { defer wg.Done() shard := &c.shards[i] shard.mx.Lock() // 复制分片内的订阅者列表,快速释放锁 subs := make([]chan<- int, len(shard.subs)) copy(subs, shard.subs) shard.mx.Unlock() for _, sub := range subs { select { case sub <- msg: case <-time.After(50 * time.Millisecond): go c.Unsubscribe(sub) } } }(i) } wg.Wait() } } func (c *Channel) Subscribe(sub chan<- int) { shard := c.getShard(sub) shard.mx.Lock() shard.subs = append(shard.subs, sub) shard.mx.Unlock() } func (c *Channel) Unsubscribe(sub chan<- int) { shard := c.getShard(sub) shard.mx.Lock() updated := make([]chan<- int, 0, len(shard.subs)) for _, s := range shard.subs { if s != sub { updated = append(updated, s) } } shard.subs = updated shard.mx.Unlock() }
优缺点:
- ✅ 锁粒度极小,每个分片的锁竞争概率只有原来的1/16(取决于分片数)
- ✅ 避免了大切片复制的开销,增删操作只针对单个分片的小切片
- ❌ 实现稍复杂,需要处理分片的哈希和并行发送的同步
方案3:结合订阅者生命周期管理减少无效操作
不管用哪种方案,都要注意及时清理无效的订阅者,避免无效的消息发送和内存泄漏:
- 当用户关闭WebSocket时,主动调用
Unsubscribe方法移除对应的channel - 在消息发送超时后,自动移除无响应的订阅者
- 给订阅者的channel设置合理的缓冲区,减少发送阻塞的概率
总结
如果你的订阅者数量在1w级别以内,方案1的atomic.Value完全够用,实现简单且性能足够;如果订阅者数量超过10w或者并发量极高,方案2的分段锁是更优的选择。另外,你之前考虑的“每个订阅者用带超时的goroutine处理发送”可以和上面的方案结合,进一步减少分发过程中的阻塞,让锁的持有时间降到最短。
内容的提问来源于stack exchange,提问作者user17644273
相关产品推荐
相关产品推荐

