因为sync.map不支持对value(即[]chan)的原子操作,load后append再store会被并发覆盖,导致订阅丢失或panic;正确做法是每个topic配独立*sync.rwmutex+[]chan,sync.map仅存该结构体指针,增删查均通过锁保护切片。

用 sync.Map 存 map[string][]chan interface{} 为什么总 panic?
因为 sync.Map 不支持对切片值的原子追加。你 Load 出一个 []chan interface{},append 后再 Store,中间可能被其他 goroutine 覆盖——结果就是订阅漏掉、切片长度错乱、甚至 nil pointer dereference。
- 正确做法:每个 topic 对应一个独立的
*sync.RWMutex+[]chan interface{}结构,sync.Map只存 topic → 这个结构体指针 - Subscribe 时先
Load,没命中就新建结构体并Store;再加写锁,append新 channel - Unsubscribe 必须加写锁遍历并删除对应 channel,不能只靠
sync.Map.Delete - 发布时用读锁遍历,但注意:遍历中不能修改切片,否则可能触发 copy-on-write 导致删不干净
select { case ch 真的够用吗?
够用,但只在 channel 有缓冲且容量合理时成立。如果 ch 是无缓冲或已满,default 会直接跳过,消息静默丢弃——这在日志通知类场景可接受,在订单状态同步类场景就是 bug。
- 带缓冲 channel 至少设为
make(chan interface{}, 16),太小易阻塞,太大吃内存 - 非阻塞发送后建议记录丢弃数(比如用
atomic.AddInt64(&dropped, 1)),便于监控背压 - 若业务不可丢,就得换方案:要么加重试队列(本地用
time.AfterFunc回推),要么直接上redis.PubSubConn - 别在
select里写case ——发布逻辑本身不该被取消,那是订阅者的事
订阅者退出后,goroutine 和 channel 为啥还在涨?
因为你没绑定生命周期控制。HTTP handler 启动一个 go func() { for range ch { ... } }(),handler 返回了,goroutine 还卡在 range 里等永远不会来的消息,channel 也一直挂在 sync.Map 里。
- Subscribe 方法必须返回
func()取消函数,内部做两件事:关闭 channel + 从 topic 结构体中删掉它 - 接收循环必须是
for { select { case v, ok := ,不能只靠 <code>range - 别把
http.Request.Context()直接传给长期运行的订阅逻辑——它的生命周期太短,容易误关 - 测试时用
sync.WaitGroup等待 goroutine 退出,比time.Sleep更可靠
要不要用 github.com/ThreeDotsLabs/watermill?
不要。它是为 Kafka / RabbitMQ 设计的重型框架,本地内存 Pub/Sub 场景下,启动慢、配置绕、日志爆炸,还强制你实现 Handler 接口和消息序列化逻辑。
- 纯内存事件分发,手写够用:200 行以内能跑通 Subscribe/Publish/Unsubscribe
- 真要跨进程或持久化,直接用
github.com/go-redis/redis/v9的Publish/Subscribe,它自带重连、断线补偿、多 subscriber 隔离 - Redis Pub/Sub 不保证可达,适合实时通知;需要 exactly-once 就上 NATS JetStream,但得额外运维
- 别用 SQLite 或文件模拟 broker——并发写入锁、崩溃恢复、消息去重全是坑
真正难的不是怎么发消息,而是谁来关 channel、谁来清理 map 里的 dead 订阅、下游处理失败时要不要重试——这些没有标准答案,得看你的业务能容忍几秒延迟、几次丢失、多少内存泄漏。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











