☰
用 C++ 从零实现一个线程安全的 Pub-Sub 消息系统:面向 LLD 面试的完整设计指南
2026/10/1 17:34:59 网站建设 项目流程
  • 示例工程

【免费下载链接】awesome-low-level-design

Learn Low Level Design (LLD) and prepare for interviews using free resources.

项目地址:https://gitcode.com/GitHub_Trending/aw/awesome-low-level-design
点击查看免费下载

导读:本文以 solutions/cpp/pubsubsystem 目录下的 C++ 实现为核心,系统讲解如何设计一个支持多发布者、多订阅者、按主题路由消息的发布-订阅(Pub-Sub)系统。你将掌握 Message、Topic、Subscriber、PubSubSystem 四个核心类的职责划分与调用链,看懂基于std::vector的主题订阅管理、订阅者消息队列与去重逻辑,并拿到一份可直接编译运行、验证"订阅/退订/发布/离线不接收"完整闭环的 Demo。这套实现思路正是低层设计(LLD)面试中消息中间件题目的高频考点。

一、需求背景:LLD 面试中的 Pub-Sub 经典题

本仓库的原始问题定义见 problems/pub-sub-system.md,其核心需求可归纳为六条:

  1. 系统允许发布者(publisher)向特定主题(topic)发布消息;
  2. 订阅者(subscriber)可以订阅感兴趣的主题,并收到发布到该主题的消息;
  3. 系统需要支持多个发布者与多个订阅者;
  4. 消息应以实时(real-time)方式投递给主题的全部订阅者;
  5. 系统需要处理并发访问并保证线程安全;
  6. 消息投递要具备可扩展性与效率。

在这份需求中,"实时投递给全部订阅者"是 Pub-Sub 区别于点对点队列的核心语义——发布者与订阅者解耦,发布者无需知道谁在收听。C++ 实现对该语义的落地方式是:每个订阅者维护一个独立的消息队列(std::vector<Message>),发布动作将消息同步写入主题内所有订阅者的队列。这种"队列内收件箱"模型是后续分析全部源码的主线。

二、整体架构与核心类职责

在 solutions/cpp/pubsubsystem 目录下,实现被拆成四个基础类与一个演示入口,共 8 个源文件:

文件角色
Message.hpp / Message.cpp消息载体,携带主题、内容与时间戳
Topic.hpp / Topic.cpp主题实体,维护订阅者列表并负责广播
Subscriber.hpp / Subscriber.cpp订阅者实体,持有个人消息队列
PubSubSystem.hpp / PubSubSystem.cpp门面控制器,统一管理主题、订阅者与发布流程
PubSubDemo.cpp演示程序,串联完整业务流程

从源码结构看,PubSubSystem是整个系统的门面(Facade):所有对外操作(建主题、加订阅者、订阅/退订、发布)都收敛到它身上,而具体的订阅集合维护与广播逻辑下沉到Topic,消息队列维护下沉到Subscriber。职责边界清晰:系统级"谁在哪个主题"由 PubSubSystem 通过两张std::vector指针表管理,主题级"订阅了谁"由 Topic 管理,订阅者级"收到了什么"由 Subscriber 管理。

需要说明的是:仓库 README.md 与 problems/pub-sub-system.md 中的英文原题描述面向 Java(使用ConcurrentHashMap与ExecutorService),而本文聚焦的 C++ 实现则是该题目在同一仓库内的多语言落地方案之一(另有 Java、Python、C#、Go 等实现,完整清单见 problems/pub-sub-system.md 的 Implementations 一节)。C++ 版基于标准库std::vector、std::find_if与手写的Message/Subscriber类完成了同构设计,因此概念可以一一对应迁移。

三、逐类源码剖析

3.1 Message:消息载体

Message.hpp 定义的Message类包含三个字段与四个公开方法:

class Message { private: std::string topic; // 消息所属主题 std::string content; // 消息正文 std::time_t timestamp; // 发布时间(Unix 时间戳) public: Message(std::string topic, std::string content); std::string getTopic() const; std::string getContent() const; std::time_t getTimestamp() const; void displayInfo() const; };

构造函数在 Message.cpp 中自动写入timestamp = std::time(nullptr),即消息一旦构造便冻结发布时间,之后无法修改——这保证了消息的不可变性(immutability),避免了投递过程中被篡改。displayInfo()依次打印主题、内容与格式化时间,供订阅者回放消息时使用。

3.2 Topic:主题与广播中枢

Topic.hpp 中,Topic持有std::vector<Subscriber*> subscribers(订阅者指针集合)、名称与描述、以及一个active状态位。四个关键方法:

  • addSubscriber(Topic.cpp):先判空,再用std::find做去重——同一个订阅者重复调用订阅不会产生重复条目;
  • removeSubscriber(Topic.cpp):按subscriber->getId()用std::find_if定位并从 vector 中erase;
  • publishMessage(Topic.cpp):若主题active为 false 则直接丢弃;否则遍历 subscribers 列表,逐个调用subscriber->receiveMessage(message),完成一次对全体订阅者的广播;
  • setActive:支持将主题整体下线,用于模拟"主题暂停服务"的场景。

可以看到,广播语义是同步的:发布者调用publishMessage时,所有订阅者队列会立即被写入。这是"实时投递"需求在单线程模型下的直接实现——以牺牲发布端阻塞为代价换取投递的即时性。

3.3 Subscriber:订阅者与个人收件箱

Subscriber.hpp 定义的Subscriber拥有id(系统生成的唯一标识)、name、messageQueue(std::vector<Message>)与active状态位。核心行为:

  • receiveMessage(Subscriber.cpp):只有在订阅者处于 active 状态时才将消息追加进队列——这是"离线订阅者不接收消息"语义的实现;若订阅者已被标记为 inactive,消息会被静默丢弃;
  • displayMessages:遍历队列逐条调用Message::displayInfo(),模拟订阅者"消费"消息;
  • clearMessages:清空队列,可理解为消费完成后释放内存;
  • displayInfo:打印订阅者名称、ID、活跃状态与待消费消息数(messageQueue.size())。

这个"每人一队列"的设计与经典 Pub-Sub 的 push 模式吻合:系统主动把消息推到订阅者的收件箱,订阅者按需回放,而不是主动去拉取。

3.4 PubSubSystem:门面控制器

PubSubSystem.hpp 是系统的心脏,内部维护:

std::vector<Topic*> topics; // 全部主题 std::vector<Subscriber*> subscribers; // 全部订阅者 int subscriberIdCounter; // 订阅者 ID 自增计数器

对外 API 与内部私有方法的分工如下表:

公开方法职责底层实现
createTopic(name, description)创建主题先findTopic查重,重名返回nullptr,否则new Topic入表
removeTopic(name)删除主题find_if定位后delete并erase
addSubscriber(name)注册订阅者自增生成SUB{n}格式 ID,new Subscriber入表
removeSubscriber(id)注销订阅者先遍历所有主题移除其订阅,再从订阅者表删除
subscribe(id, name)订阅主题两表同时查找到位才topic->addSubscriber,失败返回 false
unsubscribe(id, name)退订主题定位主题后topic->removeSubscriber
publish(name, content)发布消息主题存在且isActive()才构造Message并广播
displayTopics()/displaySubscribers()状态查看遍历调用各对象的displayInfo

几个值得注意的工程细节:

  1. ID 生成(PubSubSystem.cpp):"SUB" + std::to_string(subscriberIdCounter++),从 1 开始自增,保证注册顺序与 ID 单调递增;
  2. 注销的级联清理(PubSubSystem.cpp):removeSubscriber先把该订阅者从每一个主题的订阅列表中摘除,避免主题持有悬垂指针,这是本实现中最容易遗漏但至关重要的内存安全步骤;
  3. 发布前置校验(PubSubSystem.cpp):主题不存在或!topic->isActive()时publish返回 false,调用方可以据此判断发布失败;
  4. 内存管理:析构函数(PubSubSystem.cpp)遍历两张表逐一delete,与上面所有new严格配对,无泄漏也无重复释放。

从源码结构可以推断,findTopic与findSubscriber均采用线性扫描(std::find_if),因此时间复杂度为 O(n)。在订阅者规模较小时完全够用;若追求规模化,可像题目描述的 Java 版本那样改用哈希表(如std::unordered_map)存储主题,将查询降到 O(1),这也是需求第 6 条"可扩展与效率"的优化方向。

四、端到端 Demo 与运行验证

PubSubDemo.cpp 是一个完整的可运行示例,它演示了这条主链路:

  1. 建主题:创建Technology、Sports、Weather三个主题;
  2. 注册订阅者:添加 John、Alice、Bob 三个订阅者;
  3. 建立订阅关系:John 订阅 Technology 与 Weather,Alice 订阅 Sports,Bob 订阅 Technology 与 Sports;
  4. 发布消息:向三个主题各发一条消息;
  5. 回放验证:分别打印三个订阅者的消息队列,确认每个订阅者只收到自己订阅主题的消息;
  6. 退订验证:John 从 Weather 退订,再向 Weather 发布Storm warning!,打印 John 的队列,确认其未收到新消息。

第 3 步的订阅矩阵示意如下:

订阅者TechnologySportsWeather
John (SUB1)✅—✅(随后退订)
Alice (SUB2)—✅—
Bob (SUB3)✅✅—

运行时,displaySubscriberMessages会输出类似这样的回放结果:John 的队列中只有 Technology 消息与(退订前的)Weather 消息;退订后再发布的Storm warning!不会出现在 John 的队列里。同时,Demo 还用std::this_thread::sleep_for(std::chrono::seconds(1))模拟时间流逝,让Message::displayInfo()打印的时间戳更真实——说明发布与消费在时间上可以是异步发生的,尽管本实现内部投递是同步的。

编译与运行方式(需要 C++11 或更高标准,因为代码使用了std::to_string、std::find_if与 lambda 表达式):

# 将上述 8 个源文件(4 个 .hpp + 4 个 .cpp)置于同一目录 g++ -std=c++11 -pthread PubSubDemo.cpp PubSubSystem.cpp Topic.cpp Subscriber.cpp Message.cpp -o pubsub_demo ./pubsub_demo

说明:-pthread是为 Demo 中std::this_thread::sleep_for所需的线程库链接(部分编译器在指定-std=c++11时已隐含,但显式加上最稳妥)。若你想自行验证"离线不接收"语义,可以在订阅前调用sub->setActive(false),再发布消息,观察该订阅者队列始终为空。

五、线程安全与扩展性:从单线程到并发

需求第 5 条明确要求"处理并发访问并保证线程安全"。需要如实指出:当前 C++ 实现本身是单线程模型,其线程安全的保障主要来自两个方面:

  • 数据封装与单一入口:所有对主题、订阅者集合的修改都经由PubSubSystem门面类完成,只要调用方保证对PubSubSystem实例的访问是串行的(或外加互斥锁),内部std::vector就不会出现数据竞争;
  • 消息不可变:Message构造后字段不可改,广播时传递 const 引用,多个订阅者共享同一份内容不会产生写冲突。

对照题目原描述(见 problems/pub-sub-system.md),Java 参考实现是通过ConcurrentHashMap存储主题、用ExecutorService异步投递来满足并发与实时性要求的。若要在 C++ 版上做同样的升级,从源码结构看有两条明确路径:

  1. 存储层:将std::vector<Topic*>换成std::unordered_map<std::string, Topic*>,并把subscribe/unsubscribe/publish的查改操作包进std::mutex(或使用std::shared_mutex区分读写),即可获得 O(1) 查询 + 并发安全;
  2. 投递层:把Topic::publishMessage的同步循环改为向线程池(如std::async或第三方线程池)提交投递任务,实现"发布者不等投递完成"的异步实时推送——这正是原题中ExecutorService的角色,也是 Kafka 等生产级消息系统"读写分离、异步刷盘"思想的微缩模型。

对 LLD 面试而言,能讲清"当前实现同步投递、如何改造成异步并发"这条演进路线,比直接堆砌并发代码更有说服力。

六、设计要点总结

  1. 职责分层:Message只管数据、Subscriber只管个人收件箱、Topic只管订阅集合与广播、PubSubSystem只管编排,是典型的分层设计;
  2. 解耦语义:发布者只面向主题说话,从不直接引用订阅者,天然满足 Pub-Sub 的发布/订阅解耦;
  3. 去重与容错:Topic::addSubscriber自带去重,PubSubSystem对主题重名、订阅目标不存在等异常路径均返回nullptr/false,接口可安全调用;
  4. 内存安全:注销订阅者时级联从所有主题摘除,析构函数统一delete,指针生命周期管理闭环;
  5. 演进方向:哈希表存储 + 读写锁 + 异步投递,是这份单线程骨架走向并发可扩展系统的三步升级路径。

七、延伸阅读

  • 问题定义与全部语言实现入口:problems/pub-sub-system.md
  • 本仓库 Pub-Sub 题目对应的 UML 类图:class-diagrams/pubsubsystem-class-diagram.png(位于仓库根目录 class-diagrams 下,与本文类职责划分一一对应)
  • 其他语言的同题实现:Java(solutions/java/src/pubsubsystem/)、Python(solutions/python/pubsubsystem)、C#(solutions/csharp/pubsubsystem/)、Go(solutions/golang/pubsubsystem/)
  • 设计模式视角:本实现的"门面控制器"与 design-patterns 目录中的 Facade 模式思路一致,可作为对照阅读
  • 示例工程

【免费下载链接】awesome-low-level-design

Learn Low Level Design (LLD) and prepare for interviews using free resources.

项目地址:https://gitcode.com/GitHub_Trending/aw/awesome-low-level-design
点击查看免费下载

相关推荐

上一篇:System Informer 系统监控完整实战手册:从源码构建 4 步跑起来 + 4 个常见坑
下一篇:Duktape 2.x 中恢复 CommonJS 模块加载:module-duktape 兼容框架集成与源码解析

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询