用 sync.map + chan interface{} 实现线程安全本地事件总线:sync.map 存 topic→sublist 映射,sublist 内部用 sync.rwmutex 保护切片;每个订阅者配带缓冲 channel(如 make(chan interface{}, 10));publish 异步执行,对每个 ch 使用 select 非阻塞投递。

用 sync.Map + chan interface{} 实现线程安全的本地事件总线
进程内轻量级事件解耦,不用引入 Kafka 或 NATS 也能跑起来,但必须避开并发写 panic 和阻塞发布。
-
sync.Map存 topic →*subList映射,*subList内部用sync.RWMutex保护切片,不能只锁外层——否则并发append仍会 crash - 每个订阅者配一个带缓冲的
chan interface{}(比如make(chan interface{}, 10)),别用无缓冲 channel:一卡全卡,Publish直接 hang 住 HTTP 请求 -
Publish必须异步执行,且对每个ch做非阻塞投递:select { case ch ,否则慢订阅者拖垮全局 - 主题名加业务域前缀,比如
"user.created"而不是"created",后期对接消息中间件或做监控切片时少一半联调成本
订阅回调里为什么 goroutine 会泄漏?
不是没 go 就安全,而是没控制生命周期就危险。常见于 HTTP 调用、DB 查询这类阻塞操作没设超时。
- 回调函数内部必须自行处理
ctx.Done()或显式 timeout,比如http.Client{Timeout: 5 * time.Second},不能指望Publish统一加 timeout——它会误杀已执行成功的回调 - 别在
Subscribe里直接启动消费 goroutine;Broker 只负责分发,谁订阅谁自己开 goroutine 读 channel,否则无法按需 cancel - 若用闭包方式订阅(
Subscribe("log", func(v interface{}) { ... })),该闭包不可序列化,没法跨服务或对接 Redis/Kafka,仅限单体场景
用 go-redis 做 Pub/Sub 时 ReceiveTimeout 为什么总报 i/o timeout?
毫秒级 ReceiveTimeout 是典型误用,Redis Pub/Sub 连接专用于事件流,短轮询等于持续施压。
- 首选
pubsub.Receive()阻塞等待,它能正确处理*redis.Message、*redis.Subscription和*redis.Pong三类事件;忽略Pong会导致心跳失联,连接异常断开 -
ReceiveTimeout仅在多路复用或强制控制 goroutine 生命周期时才用,且建议设为 1~5 分钟,不是 100ms -
pubsub实例不能defer Close()在初始化函数里——它要长期存活;应在服务 shutdown 时显式调用Close()
取消订阅为什么总是失效?
不是逻辑错,是 Go 函数比较机制导致的:闭包、方法值、匿名函数即使内容相同,== 也返回 false。
-
Unsubscribe必须依赖唯一 id(比如注册时返回的stringtoken),而不是传入原函数值去比对 - 封装回调时用结构体:
type subscriber struct { id string; fn func(interface{}) },sync.Map存id → subscriber,取消时查 id 删除 - 动态订阅极强(如每秒上千次 Subscribe/Unsubscribe)时,
sync.Map的哈希竞争反而拖慢,应改用分段锁或 runtime.SetFinalizer 配合手动清理
真正难的不是搭起 Pub/Sub,是让每个订阅者既不拖慢别人,又不自己泄漏,还不被取消掉——这些细节不在接口文档里,但在上线后第一波流量里全会暴露。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











