ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

深入剖析发布-订阅(Pub-Sub)系统设计:从需求到 Go 源码级实现

深入剖析发布-订阅(Pub-Sub)系统设计:从需求到 Go 源码级实现 示例工程【免费下载链接】awesome-low-level-designLearn Low Level Design (LLD) and prepare for interviews using free resources.项目地址https://gitcode.com/GitHub_Trending/aw/awesome-low-level-design点击查看免费下载导读本文以 awesome-low-level-design 仓库中 solutions/golang/pubsubsystem 的实现为主线完整讲解发布-订阅Publisher-Subscriber系统的核心需求、类/接口设计以及 Go 语言落地实现。你将掌握 Topic、Subscriber、Publisher 的角色划分理解如何用sync.RWMutex保障并发安全并能运行仓库自带的 pubsub_system_demo.go 演示多发布者、多订阅者的实时消息投递场景。本文适用于准备系统设计/低层设计LLD面试、或需要快速搭建进程内消息总线的开发者。一、系统需求Requirements关联文档 README.md 定义了本系统需要满足的 6 条核心需求它们是后续所有类设计的出发点面向主题发布系统应允许发布者Publisher将消息发布到指定的主题Topic。按主题订阅订阅者Subscriber可以订阅感兴趣的主题并接收发布到这些主题上的消息。多对多支持系统应支持多个发布者和多个订阅者。实时投递消息应实时投递给主题的所有订阅者。并发安全系统应处理并发访问并确保线程安全。可扩展与高效系统在消息投递方面应具备可扩展性和高效性。这 6 条需求定义了一个进程内、基于主题解耦的广播模型发布者与订阅者互不感知只通过 Topic 这一中间媒介建立联系。二、核心类、接口与枚举设计关联文档给出了 7 个核心设计元素逐一展开如下Message消息表示一条可被发布、可被订阅者接收的消息内容为消息正文。对应 Go 实现见 message.go。Topic主题消息发布的目标。维护一组订阅者集合提供添加/移除订阅者以及向所有订阅者发布消息的方法。对应 topic.go。Subscriber订阅者接口定义订阅者的契约声明onMessage方法该方法在订阅者收到消息时被调用。对应 subscriber.go。PrintSubscriber打印订阅者Subscriber接口的具体实现接收消息并打印到控制台。对应 print_subscriber.go。Publisher发布者向指定主题发布消息。对应 publisher.go。PubSubSystem系统主类管理主题、订阅者与消息发布。按文档描述它使用ConcurrentHashMap存储主题、用ExecutorService处理并发消息发布——这是 Java 版 PubSubService.java 的设计对应仓库内 UML 类图 pubsubsystem-class-diagram.png而 Go 版本采用“Topic 自持读写锁 按主题同步广播”的等价方案将并发控制下沉到每个 Topic。PubSubDemo演示类通过创建主题、订阅者、发布者并发布消息来演示系统用法。Go 版对应 pubsub_system_demo.go 中的Run()入口。三、Go 源码级实现剖析3.1 Message轻量消息载体type Message struct { Content string } func NewMessage(content string) *Message { return Message{Content: content} }message.go 仅保留Content字段并通过NewMessage构造器统一创建。实际业务场景中可在此基础上扩展Timestamp、Topic、Headers等元数据Java 类图中的Message即携带timestamp: Instant与payload: String两个字段见 pubsubsystem-class-diagram.png。3.2 Subscriber 接口与 PrintSubscriber 实现type Subscriber interface { OnMessage(message *Message) }subscriber.go 定义了订阅者的唯一契约OnMessage。得益于 Go 接口的鸭子类型任何实现该方法的类型都能成为订阅者。仓库提供的默认实现 print_subscriber.gotype PrintSubscriber struct { Name string } func NewPrintSubscriber(name string) *PrintSubscriber { return PrintSubscriber{Name: name} } func (ps *PrintSubscriber) OnMessage(message *Message) { fmt.Printf(Subscriber %s received message: %s\n, ps.Name, message.Content) }每个订阅者通过Name区分身份OnMessage将消息打印到控制台——这也正是需求 4“消息实时投递给所有订阅者”的可观察落点。若需接入真实业务只需实现新的OnMessage逻辑如写入队列、调用下游 API。3.3 Topic订阅注册表 广播中枢type Topic struct { Name string Subscribers map[Subscriber]struct{} mu sync.RWMutex } func NewTopic(name string) *Topic { return Topic{ Name: name, Subscribers: make(map[Subscriber]struct{}), } } func (t *Topic) AddSubscriber(subscriber Subscriber) { t.mu.Lock() defer t.mu.Unlock() t.Subscribers[subscriber] struct{}{} } func (t *Topic) RemoveSubscriber(subscriber Subscriber) { t.mu.Lock() defer t.mu.Unlock() delete(t.Subscribers, subscriber) } func (t *Topic) Publish(message *Message) { t.mu.RLock() defer t.mu.RUnlock() for subscriber : range t.Subscribers { subscriber.OnMessage(message) } }topic.go 是整套系统的核心包含三个设计要点订阅集合用map[Subscriber]struct{}以接口值作键天然去重同一订阅者重复AddSubscriber不会产生重复投递struct{}作为空值占位零内存开销。读写锁sync.RWMutex保证并发安全对应需求 5AddSubscriber/RemoveSubscriber写操作加写锁LockPublish读操作加读锁RLock允许多个发布者并发广播、同时阻塞写入期间的集合变更。广播采用“读锁快照式遍历”Publish在持有 RLock 期间遍历订阅者并同步调用OnMessage保证发布瞬间的订阅集合一致性代价是投递是同步的、按调用者线程串行完成。3.4 Publisher受控发布type Publisher struct { Topics map[*Topic]struct{} } func NewPublisher() *Publisher { return Publisher{Topics: make(map[*Topic]struct{})} } func (p *Publisher) RegisterTopic(topic *Topic) { p.Topics[topic] struct{}{} } func (p *Publisher) Publish(topic *Topic, message *Message) { if _, exists : p.Topics[topic]; !exists { fmt.Printf(This publisher cant publish to topic: %s\n, topic.Name) return } topic.Publish(message) }publisher.go 引入了一个文档中未展开、但源码里明确实现的发布权限控制机制发布者必须先RegisterTopic(topic)登记主题才能调用Publish未登记的主题会被拒绝并打印提示。这一设计让“发布者只允许向授权主题发布”成为显式约束是对需求 1 的工程化加强。3.5 PubSubDemo完整的端到端演示func Run() { // Create topics topic1 : NewTopic(Topic1) topic2 : NewTopic(Topic2) // Create publishers publisher1 : NewPublisher() publisher2 : NewPublisher() // Create subscribers subscriber1 : NewPrintSubscriber(Subscriber1) subscriber2 : NewPrintSubscriber(Subscriber2) subscriber3 : NewPrintSubscriber(Subscriber3) publisher1.RegisterTopic(topic1) publisher2.RegisterTopic(topic2) // Subscribe to topics topic1.AddSubscriber(subscriber1) topic1.AddSubscriber(subscriber2) topic2.AddSubscriber(subscriber2) topic2.AddSubscriber(subscriber3) // Publish messages publisher1.Publish(topic1, NewMessage(Message1 for Topic1)) publisher1.Publish(topic1, NewMessage(Message2 for Topic1)) publisher2.Publish(topic2, NewMessage(Message1 for Topic2)) // Unsubscribe from a topic topic1.RemoveSubscriber(subscriber2) // Publish more messages publisher1.Publish(topic1, NewMessage(Message3 for Topic1)) publisher2.Publish(topic2, NewMessage(Message2 for Topic2)) }pubsub_system_demo.go 完整覆盖了文档演示类描述的所有动作且特意构造了跨主题订阅subscriber2同时订阅 Topic1 与 Topic2和中途退订两个边界场景步骤动作预期效果创建2 个 Topic、2 个 Publisher、3 个 Subscriber多发布者/多订阅者拓扑就绪订阅Topic1←{S1,S2}Topic2←{S2,S3}Subscriber2 同时订阅两个主题发布P1 发 2 条到 Topic1P2 发 1 条到 Topic2各主题订阅者分别收到消息退订RemoveSubscriber(subscriber2)Subscriber2 不再收到 Topic1 后续消息再发布P1 发 Message3P2 发 Message2验证退订生效S2 只收到 Topic2 的新消息四、运行方式Go 实现位于仓库的 solutions/golang 模块下模块声明见 go.modgo 1.23.2。运行有两种方式方式一通过统一入口运行。仓库根入口 solutions/golang/main.go 预留了所有项目的Run()调用点取消对应行的注释并执行go run .方式二直接在 pubsubsystem 包内验证。将Run()改为临时main()后执行go run pubsub_system_demo.go运行后控制台将依次输出类似Subscriber Subscriber1 received message: Message1 for Topic1的日志可用于直接验证“实时投递”“多订阅者广播”“退订后不再接收”等需求是否成立。五、并发安全与扩展性分析对照需求 5 与需求 6从源码结构可以做出如下分析并发安全Go 版以Topic.mu sync.RWMutex作为唯一同步原语覆盖了订阅集合的读写全部路径无裸露的共享可变状态Publisher.Topics在演示中仅在初始化阶段写入属于线程启动前的配置数据。这是比文档描述的 Java 方案ConcurrentHashMapExecutorService见 PubSubService.java更简洁的等价实现——Go 通过RWMutex把“读多写少”的广播场景优化为并发读。可扩展性当前Publish是同步串行广播订阅者数量增加时投递时延随之线性增长这是该实现的主要瓶颈点。如需扩展可参考文档所述 Java 方案的思路——将OnMessage调用提交到ExecutorService异步执行、或为每个 Topic 分配独立 goroutine 队列这些属于基于原文档设计意图的演进方向仓库当前 Go 实现并未包含。六、总结发布-订阅系统的本质是用 Topic 解耦“谁生产”与“谁消费”发布者只认主题订阅者只收回调。本仓库的 Go 实现用 6 个文件、约 130 行代码以sync.RWMutex精确满足了文档列出的全部 6 条需求并通过RegisterTopic的登记制与RemoveSubscriber的退订能力提供了超出需求清单的工程细节。阅读源码时建议按Message → Subscriber → Topic → Publisher → Demo的顺序推进即可完整复现一次“设计需求 → 类设计 → 并发落地 → 端到端验证”的 LLD 实战闭环。其他语言Java/C/C#/Python的对照实现可参考仓库 problems/pub-sub-system.md 中列出的各语言目录以及全局 UML 类图 pubsubsystem-class-diagram.png。赞分享示例工程【免费下载链接】awesome-low-level-designLearn Low Level Design (LLD) and prepare for interviews using free resources.项目地址https://gitcode.com/GitHub_Trending/aw/awesome-low-level-design点击查看免费下载相关推荐发布-订阅Pub-Sub系统低层设计LLD实战从需求到并发安全的 Java/Go 多语言实现发布 订阅Pub Sub系统低层设计LLD实战从需求到并发安全的 Java/Go 多语言实现 导读 本文以 problems/pub sub syst示例工程基于 C 与 Java 双实现剖析线程安全的发布-订阅Pub-Sub系统设计基于 C 与 Java 双实现剖析线程安全的发布 订阅Pub Sub系统设计 发布 订阅Pub Sub模式是解耦消息生产者与消费者的核心架构范式在实示例工程yuzu 模拟器Switch 游戏上 PC 的 30 分钟上手与调优指南yuzu 模拟器Switch 游戏上 PC 的 30 分钟上手与调优指南 yuzu 是一款用 C 编写的开源 Switch 模拟器由 Citra 开发团虚拟化桌面应用图形学上一篇XML Notepad智能编辑工作流突破XML处理效率瓶颈的全栈解决方案下一篇3个维度开源工具WarcraftHelper实现魔兽争霸3兼容性优化全指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表