
1. 项目概述为什么我们需要一个轻量级消息总线在C后端开发或者游戏引擎、嵌入式系统这类对性能和控制力有极致要求的场景里模块间的通信是个老生常谈但又避不开的难题。你肯定遇到过这种情况一个UI模块需要更新数据一个网络模块收到了新消息一个逻辑模块需要处理用户输入。如果让它们直接互相调用代码很快就会变成一团乱麻耦合度高到改一处动全身单元测试更是无从下手。这时候消息总线Message Bus或者事件系统Event System就成了救星。它的核心思想是“发布-订阅”Pub/Sub模块之间不直接认识对方它们只认识一个“总线”。想发消息的模块发布者把消息“扔”到总线上关心这类消息的模块订阅者从总线上“取”走并处理。这样一来发布者和订阅者完全解耦系统变得清晰、灵活、易于扩展。市面上成熟的库很多比如Boost.Signals2、各种游戏引擎内置的事件系统。但有时候我们就是需要一点“轻量级”的东西。所谓轻量级不是功能简陋而是指零外部依赖纯标准C实现编译即用、高性能避免动态内存分配、利用现代C特性优化、接口简洁易于集成和理解、类型安全告别void*和强制类型转换。自己动手实现一个不仅能完美契合项目特定需求更是深入理解异步编程、模板元编程和现代C设计模式的绝佳机会。今天我们就来从零开始打造一个属于我们自己的、生产环境可用的轻量级消息总线。2. 核心设计思路与架构选型在动手写代码之前得先把蓝图规划好。一个消息总线的核心无非是这几样消息Message、订阅者Subscriber、总线Bus本身以及它们之间的交互机制。2.1 消息定义告别裸数据拥抱类型安全最原始的做法是用一个枚举enum定义消息类型再用一个union或void*来传递数据。这种做法极其危险类型擦除和强制转换是运行时错误的温床。现代C给了我们更好的武器std::any和模板。std::any可以存储任何可拷贝的类型并在取出时进行类型检查。但它有一定的运行时开销。对于追求极致性能的场景模板是更优的选择。我们可以为每一种消息类型定义一个独立的类或结构体。总线内部使用模板来区分它们从而实现编译期的类型安全。我选择模板化主题Topic的方案。每个消息类型本身就是一个唯一的“主题”。订阅者订阅的是TopicMessageType发布者发布的也是具体的MessageType对象。这样类型匹配在编译期就完成了没有任何运行时类型信息RTTI的开销既安全又高效。2.2 订阅者管理函数对象的容器订阅者通常是一个可调用对象Callable Object普通函数、Lambda表达式、类的成员函数或者任何重载了operator()的函数对象仿函数。总线需要能存储和管理这些异构的可调用对象。这里的关键是类型擦除。我们需要一个统一的容器能存放不同类型的可调用对象。经典的解决方案是std::function。我们可以定义using Handler std::functionvoid(const MessageType)。一个主题对应一个std::vectorHandler列表。但是std::function会涉及一次堆内存分配对于捕获变量的Lambda。在性能敏感的循环中频繁订阅/退订可能成为瓶颈。更高级的优化是使用自定义的函数包装器结合小对象优化Small Object Optimization, SOO将小的可调用对象如无捕获的Lambda存储在栈内存中避免堆分配。为了首次实现的清晰度我们先使用std::function在后续优化环节再讨论替代方案。2.3 总线调度模式同步 vs. 异步消息如何传递这是架构的核心决策点。同步Sync发布消息时直接、立即、在当前线程调用所有订阅者的处理函数。优点是简单、时序确定、易于调试。缺点是如果某个订阅者处理时间过长会阻塞整个发布过程甚至导致死锁。异步Async发布消息时将消息和处理请求放入一个队列如std::queue由后台的一个或多个工作线程取出并执行。优点是不阻塞发布者能提高系统响应性天然适合多核环境。缺点是时序不确定、调试复杂、需要处理线程安全问题。对于“轻量级”的实现我建议核心提供同步总线但设计上为异步留出扩展接口。同步总线是基础满足了80%的使用场景。我们可以先实现一个线程安全的同步总线确保在多线程环境下发布和订阅操作是安全的。然后可以基于这个同步总线包装一个异步代理Async Proxy它内部包含一个任务队列和线程池将收到的消息转发给内部的同步总线去处理。这样结构清晰职责分离。2.4 总体架构图概念----------------------- | AsyncMessageBus | (可选扩展层) | (线程池 任务队列) | ---------------------- | 转发消息 ----------v------------ | SyncMessageBus | (核心层) | (线程安全的订阅者映射) | ---------------------- | 直接调用 ----------v------------ ----------------- | Subscriber List |---| std::function | | (Per Message Type) | | Handler 1 | ----------------------- ----------------- | Handler 2 | -----------------我们的实现将聚焦于核心的SyncMessageBus并确保其线程安全。3. 核心实现一步步构建线程安全的消息总线接下来我们进入具体的代码实现环节。我会分模块讲解并附上完整的代码片段。3.1 基础组件主题Topic与处理器Handler首先我们需要一个唯一的标识符来区分不同的消息类型。虽然可以用std::type_index来自typeinfo但它依赖RTTI。我们使用一个更轻量的方法利用模板和静态变量为每个类型生成一个唯一的ID。// Topic.hpp #pragma once #include cstdint namespace LightBus { namespace detail { // 内部使用的类型ID生成器 class TypeIdGenerator { static std::uintptr_t counter; public: templatetypename T static std::uintptr_t id() { static char dummy; return reinterpret_caststd::uintptr_t(dummy); // 利用函数静态变量地址的唯一性 } }; } // 主题类模板每种消息类型对应一个唯一的主题 templatetypename MessageT class Topic { public: using MessageType MessageT; static std::uintptr_t id() { return detail::TypeIdGenerator::idMessageT(); } }; }这里使用了一个小技巧模板函数idMessageT()内部的静态变量dummy对于每个不同的MessageT其地址都是唯一的。我们将这个地址转换为整数作为类型ID。这种方法在同一个程序内是唯一且稳定的且不依赖RTTI。接下来定义处理函数类型。我们使用std::function作为起点。// Handler.hpp #pragma once #include functional #include memory namespace LightBus { // 前向声明 class IMessageHandlerWrapper; // 处理器基类指针用于类型擦除的容器 using HandlerPtr std::shared_ptrIMessageHandlerWrapper; // 消息处理器包装器接口 class IMessageHandlerWrapper { public: virtual ~IMessageHandlerWrapper() default; virtual void call(const void* message) 0; // 通过void*传递消息内部进行转换 }; // 具体的消息处理器包装器模板 templatetypename MessageT class MessageHandlerWrapper : public IMessageHandlerWrapper { using Handler std::functionvoid(const MessageT); Handler m_handler; public: explicit MessageHandlerWrapper(Handler handler) : m_handler(std::move(handler)) {} void call(const void* message) override { // 将void*安全地转换回具体的消息类型指针并调用处理器 m_handler(*static_castconst MessageT*(message)); } }; }IMessageHandlerWrapper接口实现了类型擦除让我们的容器std::vectorHandlerPtr能够存储任意消息类型的处理器。MessageHandlerWrapper是模板类它保存了具体的std::function并在call方法中完成类型转换和调用。3.2 核心总线实现SyncMessageBus这是最核心的部分。总线需要维护一个映射类型ID - 该类型消息的处理器列表。同时要保证在多线程环境下订阅修改映射和发布遍历调用操作是安全的。// SyncMessageBus.hpp #pragma once #include Topic.hpp #include Handler.hpp #include unordered_map #include vector #include mutex #include shared_mutex // C17 用于读写锁 namespace LightBus { class SyncMessageBus { private: // 存储结构类型ID - 处理器列表 using HandlerList std::vectorHandlerPtr; using HandlerMap std::unordered_mapstd::uintptr_t, HandlerList; HandlerMap m_handlers; // 使用读写锁shared_mutex允许多个线程同时读取发布消息但写订阅/退订需要独占。 mutable std::shared_mutex m_mutex; public: SyncMessageBus() default; ~SyncMessageBus() default; // 禁止拷贝和赋值 SyncMessageBus(const SyncMessageBus) delete; SyncMessageBus operator(const SyncMessageBus) delete; // 订阅函数 templatetypename MessageT void subscribe(std::functionvoid(const MessageT) handler) { auto topicId TopicMessageT::id(); auto wrapper std::make_sharedMessageHandlerWrapperMessageT(std::move(handler)); { std::unique_lock lock(m_mutex); // 写锁 m_handlers[topicId].push_back(std::move(wrapper)); } } // 一个更方便的订阅重载接受任何可调用对象 templatetypename MessageT, typename Callable void subscribe(Callable callable) { subscribeMessageT(std::functionvoid(const MessageT)(std::forwardCallable(callable))); } // 发布消息同步调用 templatetypename MessageT void publish(const MessageT message) { auto topicId TopicMessageT::id(); HandlerList handlersCopy; // 复制处理器列表避免在调用时持有锁 { std::shared_lock lock(m_mutex); // 读锁 auto it m_handlers.find(topicId); if (it m_handlers.end()) { return; // 没有订阅者直接返回 } // 复制一份防止在处理器中调用subscribe/unsubscribe导致死锁或迭代器失效 handlersCopy it-second; } // 在不持有锁的情况下调用处理器 for (const auto handler : handlersCopy) { handler-call(message); } } // 退订简易版通过比较function对象来移除实际应用可能需要更复杂的标识符如token // 注意此简易实现要求传入的handler对象与订阅时是同一个比较operator这对lambda可能不适用。 // 生产环境建议返回一个订阅令牌SubscriptionToken用于退订。 templatetypename MessageT void unsubscribe(const std::functionvoid(const MessageT) handler) { // 实现略涉及在vector中查找并删除特定元素。 // 更健壮的做法是让subscribe()返回一个uint64_t的token IDunsubscribe(token)通过ID删除。 } }; }关键点解析线程安全使用std::shared_mutexC17。publish读操作使用shared_lock允许多个线程同时发布消息。subscribe写操作使用unique_lock确保修改映射表时是独占的。发布时的复制在publish函数中我们先通过读锁找到处理器列表然后立即复制一份到局部变量handlersCopy中随后释放锁再遍历这个副本进行调用。这样做至关重要它避免了“调用者持有锁”的问题。如果某个订阅者的处理函数内部又调用了subscribe或unsubscribe需要写锁就会导致死锁。复制列表虽然有一定开销但保证了安全性和可重入性。退订的复杂性上面的代码只提供了简易的unsubscribe思路。实际上由于std::function的相等性比较对于Lambda可能无效一个更通用的做法是让subscribe返回一个唯一的令牌比如一个自增的ID或弱指针退订时使用这个令牌。这部分为了核心清晰暂未实现。3.3 使用示例与测试让我们看看这个总线如何用在实际代码中。// main.cpp #include SyncMessageBus.hpp #include iostream #include string // 定义几种消息类型 struct UserLoginMessage { std::string username; int64_t timestamp; }; struct SystemAlertMessage { int level; // 1: INFO, 2: WARN, 3: ERROR std::string content; }; struct DataUpdateMessage { int dataId; double newValue; }; int main() { LightBus::SyncMessageBus bus; // 订阅UserLoginMessage bus.subscribeUserLoginMessage([](const UserLoginMessage msg) { std::cout [登录日志] 用户 msg.username 于 msg.timestamp 登录系统。 std::endl; }); // 订阅SystemAlertMessage使用函数对象 auto alertHandler [](const SystemAlertMessage msg) { const char* levelStr[] {, INFO, WARN, ERROR}; std::cout [ levelStr[msg.level] ] msg.content std::endl; }; bus.subscribeSystemAlertMessage(alertHandler); // 同一个消息类型可以有多个订阅者 bus.subscribeSystemAlertMessage([](const SystemAlertMessage msg) { if (msg.level 3) { std::cout !!! 发生严重错误请立即检查 !!! std::endl; } }); // 发布消息 bus.publish(UserLoginMessage{Alice, 1712345678}); bus.publish(SystemAlertMessage{2, 磁盘使用率超过85%}); bus.publish(SystemAlertMessage{3, 数据库连接失败}); bus.publish(DataUpdateMessage{1001, 36.5}); // 这个消息没有订阅者什么也不会发生 // 多线程测试简略 std::thread publisher([bus]() { for (int i 0; i 5; i) { bus.publish(SystemAlertMessage{1, 线程发布测试消息 std::to_string(i)}); std::this_thread::sleep_for(std::chrono::milliseconds(10)); } }); std::thread subscriber([bus]() { // 在另一个线程中动态订阅 bus.subscribeDataUpdateMessage([](const DataUpdateMessage msg) { std::cout 线程内订阅收到数据更新: ID msg.dataId , Value msg.newValue std::endl; }); }); publisher.join(); subscriber.join(); return 0; }这个例子展示了基本用法定义消息结构体、订阅、发布。代码是类型安全的如果你尝试bus.publish(DataUpdateMessage{})到一个订阅了UserLoginMessage的处理器编译器会直接报错。4. 高级特性与性能优化一个基础可用的总线已经完成了。但要用于生产环境我们还需要考虑更多。4.1 订阅令牌Subscription Token与资源管理当前的实现订阅者对象std::function的生命周期由总线内部的shared_ptr管理。如果订阅者是一个捕获了this指针的Lambda例如一个类成员函数适配器而对象在订阅后析构了那么总线在发布消息时调用这个处理器就会导致未定义行为悬空指针。解决方案使用弱引用weak subscription或订阅令牌。让subscribe返回一个Subscription对象该对象析构时自动从总线退订。这通常通过RAIIResource Acquisition Is Initialization模式实现。class Subscription { std::weak_ptrvoid m_token; // 令牌用于标识订阅 std::functionvoid() m_unsubscribeFunc; // 退订时执行的函数 public: Subscription() default; templatetypename TokenT, typename UnsubFunc Subscription(std::shared_ptrTokenT token, UnsubFunc func) : m_token(token), m_unsubscribeFunc(std::forwardUnsubFunc(func)) {} ~Subscription() { if (auto token m_token.lock()) { m_unsubscribeFunc(); } } // 禁止拷贝允许移动 Subscription(const Subscription) delete; Subscription operator(const Subscription) delete; Subscription(Subscription) default; Subscription operator(Subscription) default; }; // 在SyncMessageBus中修改subscribe返回Subscription templatetypename MessageT, typename Callable Subscription subscribe(Callable callable) { auto topicId TopicMessageT::id(); using WrapperT MessageHandlerWrapperMessageT; auto wrapper std::make_sharedWrapperT(std::functionvoid(const MessageT)(std::forwardCallable(callable))); // 为这个订阅生成一个唯一的令牌 auto token std::make_sharedint(0); // 实际可以用更轻量的结构 { std::unique_lock lock(m_mutex); // 存储时同时保存wrapper和token的弱引用 SubscriberInfo info{wrapper, token}; m_handlers[topicId].push_back(info); } // 返回Subscription对象其析构函数会执行退订逻辑 return Subscription(token, [this, topicId, token]() { std::unique_lock lock(m_mutex); auto list m_handlers[topicId]; list.erase(std::remove_if(list.begin(), list.end(), [token](const SubscriberInfo info) { return info.token.lock() token; }), list.end()); }); }这样用户可以将返回的Subscription对象作为类成员变量。当类对象析构时Subscription也随之析构自动触发退订完美解决了资源泄漏问题。4.2 性能优化替换std::functionstd::function由于类型擦除和潜在的堆分配在性能要求极高的场景如游戏每帧发布大量事件可能成为瓶颈。我们可以实现一个自定义的FunctionWrapper利用小对象优化SOO。基本思路是定义一个固定大小的缓冲区例如16或32字节如果可调用对象的大小小于缓冲区则将其直接存储在缓冲区中placement new否则再退回到堆分配。这类似于std::function的一些实现但我们可以针对我们的场景特化。templatetypename Sig class LightFunction; // 简化的轻量函数包装器 templatetypename R, typename... Args class LightFunctionR(Args...) { static const size_t BufferSize 32; using InvokeFunc R(*)(void*, Args...); using DestroyFunc void(*)(void*); using CopyFunc void(*)(void*, const void*); alignas(std::max_align_t) char m_buffer[BufferSize]; InvokeFunc m_invoke nullptr; DestroyFunc m_destroy nullptr; CopyFunc m_copy nullptr; bool m_onHeap false; templatetypename Callable static R invokeImpl(void* data, Args... args) { return (*static_castCallable*(data))(std::forwardArgs(args)...); } // ... 类似的destroyImpl和copyImpl public: templatetypename Callable, typename std::enable_if_t!std::is_same_vstd::decay_tCallable, LightFunction LightFunction(Callable callable) { using CallableType std::decay_tCallable; if (sizeof(CallableType) BufferSize alignof(CallableType) alignof(std::max_align_t)) { new(m_buffer) CallableType(std::forwardCallable(callable)); m_onHeap false; } else { m_buffer new CallableType(std::forwardCallable(callable)); m_onHeap true; } m_invoke invokeImplCallableType; m_destroy destroyImplCallableType; m_copy ©ImplCallableType; } ~LightFunction() { if (m_destroy) m_destroy(m_buffer); } R operator()(Args... args) const { return m_invoke(const_castvoid*(static_castconst void*(m_buffer)), std::forwardArgs(args)...); } // ... 移动/拷贝构造函数等 };将总线中的std::function替换为这个LightFunction可以显著减少对无捕获Lambda等小对象的动态内存分配。但实现一个完整且异常安全的LightFunction需要更多细节以上仅为概念展示。4.3 异步消息总线扩展基于同步总线构建异步层就相对简单了。我们可以实现一个AsyncMessageBus它内部包含一个SyncMessageBus实例、一个任务队列如moodycamel::ConcurrentQueue或std::queue互斥锁和一个线程池。class AsyncMessageBus { SyncMessageBus m_syncBus; moodycamel::ConcurrentQueuestd::functionvoid() m_taskQueue; // 无锁队列更优 std::vectorstd::thread m_workers; std::atomicbool m_running{true}; void workerThread() { std::functionvoid() task; while (m_running) { if (m_taskQueue.try_dequeue(task)) { task(); } else { std::this_thread::yield(); } } } public: AsyncMessageBus(size_t threadCount std::thread::hardware_concurrency()) { for (size_t i 0; i threadCount; i) { m_workers.emplace_back(AsyncMessageBus::workerThread, this); } } ~AsyncMessageBus() { m_running false; for (auto t : m_workers) t.join(); } templatetypename MessageT void publish(const MessageT message) { // 将消息和调用打包成任务放入队列 m_taskQueue.enqueue([this, message]() { // 注意这里message需要被拷贝或移动捕获 m_syncBus.publish(message); }); } // subscribe/unsubscribe 直接转发给内部的m_syncBus需要注意线程安全 };异步总线的publish操作是非阻塞的它只是将任务推入队列。工作线程从队列中取出任务并调用内部同步总线的publish从而在后台线程中同步执行订阅者的处理函数。这里需要注意消息对象的生命周期通常需要拷贝消息如上面Lambda按值捕获或者使用std::shared_ptr来传递消息。5. 常见问题、调试技巧与实战心得在实际项目中集成和使用自研消息总线肯定会遇到一些坑。这里分享几个典型问题和解决思路。5.1 问题一消息循环与递归发布导致栈溢出场景订阅者A在处理MessageX时又发布了MessageY。而订阅者B在处理MessageY时又发布了MessageX。这就形成了递归发布如果总线是同步的并且没有防护会导致调用栈不断加深最终栈溢出。解决方案禁止在消息处理函数中发布消息通过编码规范约束但这不可靠。检测递归发布在SyncMessageBus的publish方法中使用线程局部存储TLS记录当前正在处理的消息类型栈。thread_local std::vectorstd::uintptr_t s_publishStack; templatetypename MessageT void publish(const MessageT msg) { auto topicId TopicMessageT::id(); if (std::find(s_publishStack.begin(), s_publishStack.end(), topicId) ! s_publishStack.end()) { // 检测到递归记录日志或抛出异常 throw std::runtime_error(Recursive message publishing detected!); } s_publishStack.push_back(topicId); // ... 原有的发布逻辑 ... s_publishStack.pop_back(); }使用异步总线将递归调用转化为队列任务打破直接的调用链。这是最常用也最有效的办法。5.2 问题二订阅者处理时间过长阻塞主线程场景在UI线程或游戏主循环中同步发布消息某个订阅者执行了耗时操作如文件IO、复杂计算导致界面卡顿或帧率下降。解决方案耗时操作异步化这是订阅者自身的责任。订阅者收到消息后应该将耗时任务派发到工作线程然后立即返回。使用异步消息总线如前所述AsyncMessageBus的发布操作是非阻塞的处理也在后台线程。但要注意如果订阅者最终需要更新UI需要将结果再派发回UI线程例如通过另一个消息或平台提供的线程间通信机制。5.3 问题三消息顺序与丢失场景在异步模式下消息A先于消息B发布但由于线程调度B可能先于A被处理。或者在系统关闭时队列中还有未处理的消息被丢弃。解决方案顺序性简单的异步总线无法保证跨消息类型的顺序。如果需要严格顺序可以考虑使用单消费者队列或者为需要顺序的消息类型指定同一个处理线程。消息丢失在析构AsyncMessageBus时应等待队列中所有任务处理完毕优雅关闭。可以增加一个drain()方法等待队列清空。void drain() { while (!m_taskQueue.empty()) { std::this_thread::sleep_for(std::chrono::milliseconds(1)); } } ~AsyncMessageBus() { m_running false; drain(); // 等待剩余任务处理完 for (auto t : m_workers) t.join(); }5.4 调试技巧日志追踪在总线的关键路径订阅、发布、调用处理器添加详细的日志输出记录消息类型、线程ID、时间戳。这在排查“消息为什么没收到”或“死锁”问题时非常有用。静态检查利用编译期检查。例如可以定义一个宏或模板在编译时断言某个消息类型是可拷贝或可移动的因为总线可能需要存储或传递它。性能剖析使用性能分析工具如VTune、perf查看publish和处理器调用的热点。如果发现std::function调用或锁竞争是瓶颈就该考虑我们前面提到的LightFunction和无锁队列优化了。5.5 实战心得何时该用何时不该用该用消息总线的场景模块间耦合度需要降低时。需要支持动态、灵活的插件或功能扩展时。事件驱动架构如UI系统、游戏实体间的通信。需要进行全局通知但又不希望引入直接依赖时。不该滥用消息总线的场景性能极其关键的代码路径每一次消息发布都有查找、锁、函数调用的开销。对于每秒需要调用成千上万次的内部通信直接函数调用可能更合适。需要返回值或严格时序的调用消息总线是“发后即忘”fire-and-forget的。如果需要调用结果或者操作B必须严格在操作A完成后执行消息总线会增加复杂度可以考虑std::future或回调。简单的、一对一的通信如果只有两个模块需要通信且关系稳定直接接口调用更清晰。最后记住没有银弹。这个轻量级消息总线实现提供了一个坚实、可扩展的基础。你可以根据项目需求轻松地为其添加优先级、消息过滤、跨进程通信等高级功能。核心在于理解其解耦的思想和实现中的权衡性能 vs. 安全 vs. 易用性。希望这篇长文能帮助你不仅实现了一个工具更理解了其背后的设计哲学。