应使用带缓冲的 chan *Event 并配专用分发 goroutine:定义含 Type、Data、Timestamp 的 Event 结构体,启动 goroutine 持续读取分发,发送端用 select+timeout 防阻塞。

用 chan interface{} 做事件总线容易丢消息
直接用无缓冲 chan interface{} 发送事件,一旦没有 goroutine 即时接收,send 就会阻塞;加缓冲又得预估容量,爆掉后依然 panic 或静默丢弃。这不是“异步”,是伪异步——它把背压甩给了调用方,而事件总线本该负责内部缓冲和错误反馈。
实操建议:
立即学习“go语言免费学习笔记(深入)”;
- 用带缓冲的
chan *Event(不是interface{}),显式定义Event结构体,含Type、Data和Timestamp - 启动一个专用 goroutine 持续从 channel 读取并分发,避免多个消费者竞争导致漏收
- 在发送端加超时控制:
select { case bus.ch <- &event: case <-time.After(100 * time.Millisecond): log.Warn("event dropped: timeout on bus") }
回调函数注册必须支持取消订阅
用 map[string][]func(*Event) 存回调很常见,但若不提供 Unsubscribe,goroutine 持有闭包引用会导致内存泄漏,尤其在 Web handler 或短生命周期组件中反复注册时。
实操建议:
立即学习“go语言免费学习笔记(深入)”;
- 每个订阅返回一个
func()取消函数,内部用sync.Map存id → []func(*Event),避免写锁竞争 - 回调函数签名统一为
func(*Event) error,便于上游做失败重试或日志标记 - 注册时生成唯一
id(如uuid.NewString()),不要用函数地址——func类型不可比较,且热重载后地址失效
channel + 回调混合方案要注意 goroutine 泄漏
典型错误是:每次发事件都 go callback(e),但没设 context 控制生命周期。当回调里有网络请求或 sleep,旧回调可能还在跑,而订阅者早已退出。
实操建议:
立即学习“go语言免费学习笔记(深入)”;
- 所有回调执行前检查
ctx.Err() != nil,并在订阅时绑定带 cancel 的 context - 用
errgroup.Group管理回调 goroutine,主流程退出时统一Wait() - 避免在回调里直接操作共享状态(如全局 map),改用原子操作或
sync.Mutex,否则 race detector 会报错
别忽略事件重放和顺序保证
纯内存总线默认不保序也不支持重放。比如服务启动时想补推“用户登录”事件,或两个相关事件("order_created" → "payment_succeeded")被乱序处理,业务就出问题。
实操建议:
立即学习“go语言免费学习笔记(深入)”;
- 给
Event加SeqID uint64字段,消费者按需维护本地 last-seq,跳过重复或乱序事件 - 需要重放时,把历史事件写进
*bytes.Buffer或临时chan,再注入主 bus channel —— 不要直接往原 channel 写,避免干扰实时流 - 若强依赖顺序,改用单 goroutine 分发器(即“串行化投递”),而不是并发调用多个回调
实际最难的不是实现,是界定哪些事件必须可靠投递、哪些可以丢。比如“埋点上报”可容忍丢失,但“库存扣减确认”必须至少一次——这决定了你要不要引入外部存储(如 Redis Stream)或 ACK 机制。channel 和回调只是工具,边界才是设计起点。

















