订阅发布模块事例
背景由于业务之间的解耦通常我们在一个进程中会设计多个模块模块之间的通信有时不需要消息队列但是需要一个通信的机制。示例1——消息内容使用基类#includeiostream#includestring#includevector#includefunctional#includemap#includemutex#includememory// 通用消息基类可选用于统一传递数据structMessage{virtual~Message()default;};classMessageBus{public:usingCallbackstd::functionvoid(conststd::shared_ptrMessage);// 获取单例staticMessageBusinstance(){staticMessageBus bus;returnbus;}// 订阅指定字符串主题voidsubscribe(conststd::stringtopic,Callback callback){std::lock_guardstd::mutexlock(mutex_);subscribers_[topic].push_back(std::move(callback));}// 发布指定字符串主题和消息内容voidpublish(conststd::stringtopic,conststd::shared_ptrMessagemsg){std::lock_guardstd::mutexlock(mutex_);autoitsubscribers_.find(topic);if(itsubscribers_.end())return;// 复制回调列表避免执行过程中修改容器autocallbacksit-second;// 释放锁后执行防止死锁lock.unlock();for(constautocb:callbacks){try{cb(msg);}catch(...){std::cerrError in subscriber for topic: topicstd::endl;}}}private:MessageBus(){}std::mutex mutex_;// Key: 主题名称 (String), Value: 回调列表std::mapstd::string,std::vectorCallbacksubscribers_;};// 测试示例 structOrderMsg:publicMessage{intorderId;OrderMsg(intid):orderId(id){}};structLogMsg:publicMessage{std::string content;LogMsg(conststd::stringc):content(c){}};intmain(){autobusMessageBus::instance();// 1. 订阅 order 主题bus.subscribe(order,[](conststd::shared_ptrMessagemsg){autoorderstd::dynamic_pointer_castOrderMsg(msg);if(order){std::cout[OrderModule] 收到订单: order-orderIdstd::endl;}});// 2. 订阅 log 主题bus.subscribe(log,[](conststd::shared_ptrMessagemsg){autologstd::dynamic_pointer_castLogMsg(msg);if(log){std::cout[LogModule] 收到日志: log-contentstd::endl;}});// 3. 发布消息bus.publish(order,std::make_sharedOrderMsg(1001));bus.publish(log,std::make_sharedLogMsg(System Start));return0;}示例2——消息类型模板化#includeiostream#includestring#includevector#includefunctional#includemap#includemutex#includememory#includetypeindex// 1. 内部存储辅助类 // 用于在 Map 中存储不同类型的回调并管理数据的生命周期classMessageBus{public:staticMessageBusinstance(){staticMessageBus bus;returnbus;}// 订阅指定主题和回调templatetypenameTvoidsubscribe(conststd::stringtopic,std::functionvoid(constT)callback){std::lock_guardstd::mutexlock(mutex_);// 使用 type_index 确保同一主题下不同类型的回调隔离可选若约定主题唯一类型可省略autokeystd::make_pair(topic,std::type_index(typeid(T)));// 包装回调从 void* 恢复为 Tautowrapper[callback](void*data){callback(*static_castconstT*(data));};registry_[key].push_back(wrapper);}// 发布值传递内部会拷贝一份数据templatetypenameTvoidpublish(conststd::stringtopic,constTmessage){std::lock_guardstd::mutexlock(mutex_);autokeystd::make_pair(topic,std::type_index(typeid(T)));autoitregistry_.find(key);if(itregistry_.end())return;// 【关键】拷贝消息数据确保生命周期独立于发布者// 使用 shared_ptr 管理拷贝后的数据确保所有回调执行期间数据有效autodataCopystd::make_sharedT(message);autocallbacksit-second;lock.unlock();for(constautocb:callbacks){try{// 传递拷贝数据的指针回调中转换为引用cb(dataCopy.get());}catch(...){std::cerrError in subscriber.std::endl;}}}private:MessageBus(){}std::mutex mutex_;// Key: 主题名, 类型信息, Value: 类型擦除的回调列表usingKeystd::pairstd::string,std::type_index;std::mapKey,std::vectorstd::functionvoid(void*)registry_;};// 2. 测试示例 structOrderEvent{intorderId;std::string customer;OrderEvent(intid,conststd::stringname):orderId(id),customer(name){std::coutOrderEvent Created: orderIdstd::endl;}~OrderEvent(){std::coutOrderEvent Destroyed: orderIdstd::endl;}};intmain(){autobusMessageBus::instance();// 订阅bus.subscribeOrderEvent(order,[](constOrderEvente){std::cout[Subscriber] Received Order: e.orderId, e.customerstd::endl;});{// 发布传入局部变量值传递OrderEventlocalOrder(101,Bob);std::cout--- Publishing ---std::endl;bus.publish(order,localOrder);std::cout--- Publish Call Finished ---std::endl;}// localOrder 在此处销毁但订阅者收到的拷贝依然有效return0;}