go-pubsub:一个为实时流场景设计的进程内 Pub/Sub

2026-9-20

这篇文章介绍我写的第一个开源库 go-pubsub:一个轻量的进程内 Pub/Sub,专为实时流媒体数据包这类瞬态数据流设计。它的哲学是纯粹的 fire-and-forget——零持久化、不保证送达,只追求极致的单向消息传递速度。

为什么写它

做实时流媒体相关开发时,经常需要在进程内分发数据包:一个采集协程产出帧,N 个消费协程各自订阅。这类场景的共同点是:

  • 数据转瞬即逝:晚到一帧不如直接丢掉,等下一帧
  • 消费者可能随时消失:订阅方断连后不应该拖累生产者
  • 延迟比可靠性重要得多

标准库里 map[string]chan T 手搓一个并不难,但每次都要重新处理并发安全、订阅清理、背压策略这些边角问题。用 Kafka / NATS 又太重——数据根本不需要出进程。于是我把这些需求固化成了一个库。

快速上手

go get github.com/F2077/go-pubsub

基于泛型,消息类型在编译期确定:

broker, _ := pubsub.NewBroker[string]()
publisher := pubsub.NewPublisher[string](broker)
subscriber := pubsub.NewSubscriber[string](broker)

// Subscribe with a buffered channel and a sliding 200ms timeout:
// if no publish lands within the window, ErrSubscriptionTimeout
// is fired to ErrCh.
sub, _ := subscriber.Subscribe("alerts",
    pubsub.WithChannelSize[string](pubsub.Medium),
    pubsub.WithTimeout[string](200*time.Millisecond),
)
defer sub.Close()

publisher.Publish("alerts", "CPU over 90%!")

select {
case msg := <-sub.Ch:
    fmt.Println("Received:", msg)
case err := <-sub.ErrCh:
    log.Println("Timeout:", err)
}

几个关键的设计取舍

满则丢弃,而不是阻塞。 channel 满了消息就静默丢弃,生产者永远不被慢消费者拖住。这正是实时流想要的行为,但也意味着它绝对不能用于需要可靠投递的场景。

滑动超时自动回收订阅。 WithTimeout 启动一个滑动计时器,每次成功投递都会重置;窗口内没有消息就往 ErrCh 发一次 ErrSubscriptionTimeout。订阅方异常退出但没 Close 时,资源也能自动释放。

预置的 buffer 档位。 Block / Single / Small / Medium / Large / Huge 六档(Medium = 100 是默认值),也接受任意 uint16。基准测试显示 buffer 档位对发布性能几乎无影响(104~107 ns/op),按消费者的消费能力选就行。

热路径零分配。 fan-out 时订阅者快照切片用 sync.Pool 复用,单订阅者发布是 107.7 ns/op,0 allocs/opallocs/op 是确定可复现的指标,ns/op 会随机器负载波动,所以仓库里把基准测试方法和注意事项都写明了。

什么时候用,什么时候别用

  • ✅ 实时流数据分发、游戏/直播事件、进程内低延迟消息
  • ❌ 持久化队列、需要送达保证的场景——这类需求请用专业的消息队列

项目地址:https://github.com/F2077/go-pubsub,MIT 协议,欢迎 issue 和 PR。

https://since2077.top/blog/atom.xml