不能直接用chan做广播,因Go的channel是点对点原语,一次send仅被一个recv消费,多goroutine range同一chan会导致竞态抢答和阻塞;正确做法是用sync.Map存topic→subList映射,subList内用RWMutex保护[]chan并配带缓冲通道,发布时RLock快照遍历+select非阻塞投递。

为什么不能直接用 chan 做广播?
Go 的 chan 是点对点通信原语:一次 send 只能被一个 recv 消费。如果你启动 5 个 goroutine 同时 range 同一个 chan string,结果只有其中一个能收到消息,其余全部卡死在等待状态——这不是广播,是竞态抢答。
- 常见错误现象:
for range ch在多个协程里监听同一通道,每条消息只进一个协程 - 真实场景需求:WebSocket 聊天室推消息给所有在线用户,不是“谁抢到算谁的”
- 根本限制:channel 没有“复制分发”能力,必须靠代码层补足发布-订阅语义
- 性能影响:强行复用单 channel 做伪广播,会导致消息丢失、逻辑错乱,调试时极难定位
怎么用 sync.Map + chan 实现线程安全的 Pub/Sub?
核心是为每个 topic 维护一个专属的订阅通道切片,且每个订阅者拿到自己的 chan interface{},由发布者遍历并安全投递。但注意:sync.Map 只保 key 级原子性,不保护 value(比如 []chan)的并发读写。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
- 推荐结构:
sync.Map存topic → *subList,其中*subList内部用sync.RWMutex保护[]chan interface{} - 订阅时调用
subList.add(ch),发布时遍历前先RLock(),避免遍历时切片被修改 panic - 每个订阅通道必须带缓冲(如
make(chan interface{}, 16)),否则一个慢消费者会拖垮整个Publish - 发布逻辑必须用
select { case ch 非阻塞发送,跳过写不进的通道,防止卡死 - 务必提供
Unsubscribe(topic, ch),内部做原子查找+删除,不能直接操作切片
广播时如何避免 goroutine 泄漏和连接残留?
客户端断开后,若没及时清理其 send chan 和 map 条目,会导致内存缓慢上涨、广播变慢、甚至 panic:send on closed channel。
- 注册新连接时,生成带缓冲的
msgCh := make(chan []byte, 64),避免接收方慢导致广播协程卡住 - 客户端读
msgCh的 goroutine 一旦遇到io.EOF或写错误,必须主动触发注销:删 map 条目 + 关闭msgCh - 订阅方法签名建议为
Subscribe(ctx context.Context, topic string) (func()),返回的取消函数负责清理资源 - 接收循环必须包含
case ,不能忽略上下文退出信号 - 别把
http.Request.Context()直接传给长期运行的订阅逻辑——它的生命周期太短,容易误关
要不要用 github.com/ThreeDotsLabs/watermill 这类第三方库?
watermill 是面向 Kafka / RabbitMQ 的重型消息框架,本地内存 Pub/Sub 属于杀鸡用牛刀:启动慢、配置绕、trace 日志泛滥,还强制你写 handler 接口和消息结构体。
- 纯内存场景下,
sync.Map + chan + select组合已足够支撑万级 QPS,关键在规避那几个硬伤 - 真正需要第三方库的信号是:要跨进程通知、服务重启后消息不丢、支持通配符订阅(如
logs.*)、或需持久化 - 此时直接集成
redis.PubSubConn或nats.go,比自己基于 TCP 或文件硬搞更可靠 - 复杂点在于:Redis 解决的是「跨进程通知」,不是「连接状态同步」;你仍需在每个节点上维护自己的 clients 映射

















