必须加缓冲和非阻塞写,因为go channel是同步原语,无缓冲chan在订阅者处理慢时会导致publish阻塞;带缓冲chan(如容量16)配合select+default可实现非阻塞发送,避免慢消费者拖垮全局发布。

用 sync.Map + chan 实现本地 Pub/Sub,为什么必须加缓冲和非阻塞写?
因为 Go 的 channel 是同步原语,不带缓冲的 chan interface{} 一旦订阅者处理慢或卡住,发布端就会在 ch 处永久阻塞——整个系统就挂了。
- 每个订阅者 channel 必须设缓冲,比如
make(chan interface{}, 16),避免单个慢 consumer 拖垮全局 - 发布时不能直接写:
ch ,得用 <code>select { case ch 非阻塞发送 -
default分支不是“丢消息”,而是跳过当前订阅者,保证其他 subscriber 不受影响 - 缓冲大小不是越大越好:太大(如 1024)会掩盖消费瓶颈,还可能堆积大量 stale 消息导致内存上涨
sync.Map 存订阅列表,为什么不能直接用 map[string][]chan?
sync.Map 本身不支持原子追加 slice 元素,LoadOrStore 对 []chan 类型无效——你并发调用两次 Subscribe,很可能只保留最后一个,前一个被覆盖。
- 正确做法是把每个订阅者封装成结构体,含唯一 ID、channel、cancel 函数,再存进
sync.Map - 别在
Range回调里对sync.Map做写操作(比如删掉已关闭的 channel),会 panic;应先LoadAll()拷贝快照,再遍历处理 - 如果订阅关系稳定(比如服务启动后基本不增删),用普通
map+sync.RWMutex更轻量,sync.Map的优势只在千级并发注册/注销场景
订阅者退出时,怎么防止 goroutine 和 channel 泄漏?
常见错误是让订阅者自己关 channel,结果发布端还在往已关闭的 ch 写,触发 panic: send on closed channel;或者忘了关,goroutine 卡在 select { case 永远不退出。
- 订阅函数必须返回
context.CancelFunc,由调用方控制生命周期,而不是暴露close(ch) - 发布端监听
ctx.Done(),收到信号后主动从sync.Map中移除该订阅者,而不是等它自己关 - 别在 goroutine 里
defer close(ch):万一 context 先取消,defer还是会执行,造成 double-close - 测试泄漏最简单方法:启动订阅后立刻 cancel,看 goroutine 数是否回落(用
runtime.NumGoroutine()或 pprof)
要不要用 github.com/ThreeDotsLabs/watermill?
watermill 是重型框架,面向 Kafka/RabbitMQ 场景。本地内存 Pub/Sub 用它,等于拿火箭发射玩具车——启动慢、配置绕、日志泛滥,还强制你实现 MessageMarshaler 接口。
- 如果你只是做配置热更新、状态通知、WebSocket 广播这类轻量事件分发,
sync.Map + chan足够,50 行代码搞定 - 只有当需要跨进程、持久化、ack、重试、consumer group 时,才值得引入
watermill或go-redis/v8 - 用
watermill最容易踩的坑是:忘记调router.Run(),或者ConsumerGroup名动态生成,导致每次启动都算新组、丢 offset
实际项目里,最难的从来不是“怎么写”,而是“怎么安全地关”——channel 关早了丢消息,关晚了泄漏;context 取消时机不对,会导致部分订阅者残留;sync.Map 的并发语义又不像 mutex 那么直觉。这些边界条件,往往要压测到 QPS 上千才会暴露。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











