diskqueue 本身不支持发布订阅,因其 readchan() 返回唯一读通道,所有消费者争抢同一条消息且消息消费后即删除;需用 sync.map 管理多 topic 订阅通道,配合专属 goroutine 读取并分发,但仅限单机、需防 oom 和配置调优。

diskqueue 本身不支持发布订阅(pub/sub)语义,它只是个单向、有序、基于文件的 FIFO 队列。想用它实现「一个 topic 多个订阅者各自独立消费」,必须自己补全上层逻辑——否则所有消费者会争抢同一条消息,无法做到广播。
为什么直接用 diskqueue 不能当 pub/sub 用
diskqueue.ReadChan() 返回的是**唯一读通道**,所有 goroutine 从同一个 chan []byte 里读,Go 的 channel 是“抢读”模型:谁先 谁拿到消息,其余协程阻塞或跳过。这本质是点对点(P2P),不是发布订阅。
常见错误现象:
- 启动两个消费者 goroutine,但只有其中一个能持续收到消息,另一个长期阻塞
- 消息被消费一次后就从磁盘文件中删除(
diskqueue的默认行为),其他订阅者根本无从获取 - 手动克隆
ReadChan()多次,结果 panic: “send on closed channel”,因为底层 reader 关闭后通道已关
用 sync.Map + diskqueue 搭出轻量 pub/sub
核心思路:把 diskqueue 当作底层「消息持久化层」,只负责可靠落盘和顺序读取;真正的「发布-分发」逻辑由内存中的 sync.Map 承担。
实操建议:
- 启动一个专属 goroutine 从
diskqueue.ReadChan()持续读取消息,每次读到后调用publishToTopic(msg, topic) -
sync.Map存的是map[string][]chan []byte,key 是 topic 名,value 是该 topic 下所有活跃订阅者的接收通道切片 - 每个
Subscribe(topic string)返回一个新的chan []byte,并把它 append 到对应 topic 的切片中(需用sync.RWMutex保护切片操作) - 发布时遍历该 topic 的所有通道,用
select { case ch 非阻塞发送,避免某个慢订阅者拖垮整体 - 务必提供
Unsubscribe(topic string, ch chan []byte),内部做切片删除 + 关闭通道,否则 goroutine 和内存泄漏
性能与边界要盯住这三点
diskqueue 的吞吐取决于磁盘 I/O 和配置参数,而你加的内存分发层会引入额外开销。容易被忽略但关键的点:
-
diskqueue初始化时的maxBytesPerFile和syncTimeout直接影响消息刷盘频率——设太小导致频繁小写,设太大则重启后恢复慢 - 如果订阅者处理速度远低于发布速度,
chan []byte缓冲区会堆积,最终 OOM;建议所有订阅通道都设合理 buffer(如make(chan []byte, 100)),而非 0 - 没有跨进程能力:这个方案只适用于单机多协程。若需多个服务实例共享同一 topic,
diskqueue的文件路径必须共享(比如挂载同一 NFS),且需额外加分布式锁协调读 offset,复杂度陡增——这时候该换redis.PubSubConn或nats.go
别硬扛持久化,该借力时就借力
自己基于文件实现带 topic 分发、消息重试、死信队列、通配符匹配的 pub/sub,99% 的场景都是重复造轮子。真正需要文件级持久化的轻量需求,其实只集中在两类:
- 嵌入式设备或边缘节点,没条件跑 Redis/NATS,且消息量极低(
- 作为主消息系统(如 Kafka)的本地 fallback 缓存,断网时暂存待发消息
除此之外,直接上 redis-cli pubsub 或 go-redis 的 Publish/Subscribe,5 行代码起步,运维成本归零。自己写的文件队列,debug 一次 offset 错位 或 segment 文件残留 花的时间,够你搭好三套 Redis 实例。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











