go消息总线中间件应聚焦publish/handle两处钩子:beforepublish可丢弃事件,afterhandle统一处理panic与日志;须用context透传上下文,避免全局状态污染。

Go 里没有现成的“消息总线中间件”抽象,net/http 那套包装 http.Handler 的思路不能直接搬过来——消息总线的核心是事件分发与订阅,拦截点不在请求/响应链上,而在发布(publish)和消费(handle)两个环节。
为什么不能直接套用 HTTP 中间件模式
HTTP 中间件依赖 next.ServeHTTP() 控制执行流;而消息总线中,bus.Publish("user.created", data) 是个无返回值的异步调用,没有“下一个处理器”的显式链条。强行模拟会导致:
- 订阅者注册时被中间件层层包装,但实际触发靠
bus.Emit()内部调度,中间件无法自然嵌套 - 若在
Publish入口统一拦截,就只能做前置过滤(如丢弃敏感事件),无法对单个订阅者做差异化处理(比如只给 A 订阅者加日志,B 订阅者加重试) - 多个中间件叠加后,错误传播路径混乱,
recover()容易漏捕获
真正可行的拦截点只有 publish 和 handle 两处
轻量级总线(比如基于 sync.Map + chan 实现的简单事件总线)应只暴露两个可插拔钩子:
-
BeforePublish(event string, payload interface{}) bool:返回false则直接丢弃该事件,不进入分发队列 -
AfterHandle(event string, payload interface{}, err error):无论订阅者是否 panic,都保证执行,适合日志、指标上报
注意:不要试图在 Subscribe() 时传入“带中间件的 handler”,那会让类型变复杂且破坏订阅语义。正确做法是让每个订阅者本身是干净的函数,拦截逻辑下沉到总线内部统一调度层。
如何安全地在 handle 阶段拦截 panic 并恢复
订阅者函数可能 panic,但你不希望一个订阅者崩溃导致整个总线挂掉。必须在总线内部的调用点包裹 recover():
func (b *Bus) emitToSubscriber(sub *subscriber, event string, payload interface{}) {
defer func() {
if r := recover(); r != nil {
b.afterHandle(event, payload, fmt.Errorf("panic in subscriber %p: %v", sub.fn, r))
}
}()
sub.fn(event, payload)
}
关键点:
- 不能只在
Publish()外层 recover——那只能捕获调度逻辑本身的 panic,捕不到订阅者执行时的 panic -
afterHandle必须是总线实例方法,才能访问配置的拦截器回调 - 别用
log.Printf直接打 panic 日志,要走AfterHandle钩子,方便下游替换为 Sentry 上报或降级处理
中间件参数传递必须靠闭包,而非全局配置
如果你需要某个拦截逻辑读取上下文(比如当前 trace ID),不要把 context 存进总线结构体——那会污染状态、引发并发问题。正确方式是让 Publish() 接受可选的 context.Context,并在内部透传给所有钩子:
func (b *Bus) Publish(ctx context.Context, event string, payload interface{}) {
if !b.beforePublish(ctx, event, payload) {
return
}
// ... 分发逻辑
for _, sub := range b.subscribers[event] {
go b.emitToSubscriberWithContext(ctx, sub, event, payload)
}
}
这样 BeforePublish 和 AfterHandle 都能拿到 ctx,从中提取 trace.SpanFromContext(ctx) 或自定义值,而无需共享变量。
最常被忽略的是:事件总线的“中间件”本质是策略钩子,不是函数链。想加新行为?不是往链尾 append 一个函数,而是实现并注册一个 BeforePublish 回调。混淆这两者,就会写出难以调试、panic 传播不可控的总线代码。











