最常用go原生无内置pub/sub,用sync.map存topic→[]chan interface{}、chan传消息最直接可控;需防map并发写panic、chan未关闭致内存泄漏、发布端阻塞及消息丢失。

用 sync.Map + chan 实现轻量级 Pub/Sub 为什么最常用
Go 原生没有内置 Pub/Sub,但用 sync.Map 存订阅者、chan 传消息,是最直接且可控的方式。它不依赖第三方库,适合中小规模事件分发,比如配置热更新、状态通知。
常见错误是直接用 map 并发写 panic,或在 goroutine 中忘了关闭 chan 导致内存泄漏。
-
sync.Map存topic→[]chan interface{},避免读写锁争用 - 每个订阅者用带缓冲的
chan(如make(chan interface{}, 16)),防发布端阻塞 - 订阅函数返回取消函数,内部调用
close(ch)并从sync.Map删除该chan - 发布时遍历
sync.Map的值,用select配合default避免向已满或已关的chan写入
用 github.com/ThreeDotsLabs/watermill 做分布式 Pub/Sub 要注意什么
当消息需要跨进程、持久化或对接 Kafka/RabbitMQ 时,watermill 是主流选择。但它不是“开箱即用”的简单封装,配置错一步就收不到消息。
典型现象:本地测试能 publish,但 subscriber 不触发;或者消息重复消费。
- 必须显式调用
router.Run启动,否则SubscribeHandler不生效 - Kafka 场景下,
ConsumerGroup名称不能动态生成,否则每次启动都算新组,丢失 offset - 使用
GoChannel作为本地传输时,router.AddHandler的 topic 名必须和Publish时一致,大小写敏感 - 消息结构体需实现
watermill.MessageMarshaler才能自动序列化,否则报cannot marshal message
context.Context 怎么安全注入到 Pub/Sub 生命周期里
Pub/Sub 往往涉及超时控制、取消传播、trace 上下文透传——但很多人把 context.Background() 硬编码进 publish 或 subscribe 函数,导致无法中断长耗时处理。
错误示例:go func() { handler(msg) }(),此时 handler 完全脱离原始 context。
- 订阅注册时传入
context.Context,并在每个消息处理 goroutine 中用ctx, cancel := context.WithTimeout(parentCtx, 5*time.Second) - publish 函数签名应为
Publish(ctx context.Context, topic string, msg interface{}) error,内部用ctx.Done()检查是否超时再写入chan - watermill 中,用
message.NewMessage构造消息时,可通过message.WithContext注入 traceID 或 deadline - 别在
defer cancel()后继续操作 channel,cancel 后ctx.Done()可能立即关闭,引发竞态
为什么不用 channel 直接做全局 Pub/Sub
有人试过定义一个全局 chan interface{},所有发布者往里塞,所有订阅者 range 读——这在单协程下看似可行,但实际会立刻崩。
现象包括:goroutine 泄漏、消息丢失、panic: send on closed channel、CPU 100%。
- channel 无 topic 分区,所有消息混在一起,订阅者无法过滤
- range channel 会阻塞直到 channel 关闭,而全局 channel 几乎不可能优雅 close
- 一个订阅者 panic 会导致整个 range 循环退出,其他订阅者全部失联
- 无法动态增删订阅者,也不能对不同 topic 设置不同 buffer 大小或超时策略
真正复杂的点不在怎么写逻辑,而在怎么让订阅者“自己决定什么时候退出”——这个退出信号要能被发布端感知、被中间件传递、被下游处理链正确响应。稍有疏忽,就是静默丢消息或 goroutine 积压。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











