You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

网站后端通知系统:无需互斥锁的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 06:48:31