不能裸用 sync.Map 存 chan interface{},因其仅解决 map 读写竞争,不管理 channel 生命周期、关闭状态及缓冲区满载问题;Range 是快照语义,新增订阅者收不到当前发布消息,且未关闭 channel 会导致 goroutine 泄漏。

直接用 map[string][]chan interface{} + sync.RWMutex 就能跑通,但线上崩就崩在没关 channel、没做非阻塞发送、没控制 goroutine 生命周期。
为什么不能只用 sync.Map 存 chan interface{}
很多人一上来就用 sync.Map 存订阅通道,觉得“并发安全”就万事大吉。其实它只解决 map 本身的读写竞争,掩盖了更关键的问题:
-
sync.Map.Range()是快照语义:遍历时新增的订阅者不会收到本次发布的消息——这本身合理,但如果你误以为“所有当前活跃 subscriber 都会被通知到”,就会漏事件 - 值类型是
chan interface{}时,sync.Map不管 channel 是否已关闭,也不管缓冲区是否堆满;发布方往一个没人读的 channel 里写,select { case ch会直接走default丢弃,还是得靠你显式判断和日志补位 - 真正该用
sync.Map的场景,是 topic 数量多、读远多于写的总线(比如万级 topic 的监控系统);普通业务服务用带锁的map更易 debug,panic 位置清晰
Subscribe 必须返回取消函数,且内部要 close(ch)
常见错误是 Unsubscribe 只从 map 里删掉 channel,却不关它。后果很直接:
- 消费者 goroutine 还卡在
for range ch上,因为已关闭的 channel 仍可读(返回零值),无法退出 - 如果消费者用的是
for { select { case msg := ,没检查 <code>ok,就会无限循环读零值 - goroutine 泄漏肉眼不可见,但
runtime.NumGoroutine()在反复订阅/退订后持续上涨,就是典型信号
正确做法:
立即学习“go语言免费学习笔记(深入)”;
func (eb *EventBus) Subscribe(topic string, bufSize int) chan interface{} {
ch := make(chan interface{}, bufSize)
eb.mu.Lock()
defer eb.mu.Unlock()
eb.subscribers[topic] = append(eb.subscribers[topic], ch)
return ch
}
<p>func (eb *EventBus) Unsubscribe(topic string, ch chan interface{}) {
eb.mu.Lock()
defer eb.mu.Unlock()
if subs, ok := eb.subscribers[topic]; ok {
for i, c := range subs {
if c == ch {
close(ch) // 先关
eb.subscribers[topic] = append(subs[:i], subs[i+1:]...)
break
}
}
}
}
Publish 一定要用 select { case ch 非阻塞发送
别写 ch 。哪怕加了缓冲,只要某个订阅者消费慢,整个 <code>Publish 就被拖住——发布方本不该为订阅方的延迟买单。
- 每个订阅 channel 都要单独做
select,避免一个卡死拖垮全部 - 加
default后,可以统计丢弃数或打 warn 日志,而不是静默失败 - 不要在发送路径上加
ctx.Done()判断:那会掩盖 channel 是否 ready 的事实,且增加无谓开销 - 如果消息体是结构体指针(如
*User),多个订阅者并发改它会触发 data race;传值或用sync.Pool复用更安全
主题命名和消息体不规范,后期根本没法拆服务
现在用 "user.created" 还是 "created",看起来只是少两个字。但半年后你要把用户模块拆成独立服务,问题就来了:
- 没有前缀的主题无法路由:Kafka 里没法按
user.*做 topic 级限流或 ACL - 消息体用
interface{}包map[string]interface{},等于放弃编译期校验;上线后字段拼错、类型不一致全靠日志猜 - 真正该传的是定义好的结构体:
type UserCreatedEvent struct { UserID int; Email string; Timestamp time.Time },再序列化为 JSON 或 Protobuf
跨服务的解耦,从来不在 channel 怎么写,而在 topic 和 schema 怎么定。单机总线写得再漂亮,主题一乱,整个事件流就变成黑盒。


















