C++多线程观察者模式实战:写时复制与线程安全实现

发布时间:2026/7/25 5:12:39
C++多线程观察者模式实战:写时复制与线程安全实现 1. 项目概述当观察者模式遇上多线程在C后端开发里事件驱动和异步通知是家常便饭。比如一个服务状态发生变化需要立刻通知到多个监控模块或者一个数据源更新了一堆数据处理单元等着消费。这时候观察者模式Observer Pattern就成了最自然的选择——它定义了一种一对多的依赖关系让多个观察者对象同时监听某一个主题对象主题对象状态变化时会通知所有观察者对象使它们能够自动更新。但事情一放到多线程环境味道就变了。想象一下你正在一个高并发的交易系统中行情数据主题每秒更新成千上万次几十个分析策略观察者需要实时接收数据。如果还用教科书里那种简单的vectorObserver*加一个同步的notify调用轻则数据丢失、更新顺序错乱重则直接程序崩溃死锁、竞态条件全来了。这不再是简单的“设计模式”问题而是实打实的“并发编程”挑战。所以我们今天要聊的就是如何把经典的观察者模式安全、高效地移植到C的多线程世界里。这不仅仅是给std::vector加把锁那么简单它涉及到线程安全的数据结构、通知策略的选择、生命周期的管理甚至是性能与一致性的权衡。如果你正在设计一个需要处理高频事件通知的C服务或者对如何优雅地处理多线程回调感到头疼那这篇从实战踩坑中总结出来的经验应该能给你一些直接的参考。2. 核心设计思路与线程安全挑战2.1 观察者模式基础与多线程引入的变数我们先快速回顾一下观察者模式的核心骨架。通常它包含两个主要角色Subject主题维护一个观察者列表提供注册attach和注销detach方法并在状态改变时调用通知方法notify。Observer观察者定义一个更新接口如update供主题在通知时调用。一个最简单的单线程C实现可能长这样class Observer { public: virtual ~Observer() default; virtual void update(const std::string msg) 0; }; class Subject { std::vectorObserver* observers_; public: void attach(Observer* obs) { observers_.push_back(obs); } void detach(Observer* obs) { // 需要找到并删除obs这里省略实现 } void notify(const std::string msg) { for (auto* obs : observers_) { obs-update(msg); // 同步调用阻塞直到所有观察者处理完 } } };一旦进入多线程环境这个简单的模型会面临几个致命问题数据竞争Data Raceattach、detach和notify可能被多个线程同时调用。比如线程A正在遍历observers_进行通知线程B却同时删除了其中一个观察者detach这会导致迭代器失效引发未定义行为通常是崩溃。死锁Deadlock如果观察者的update方法内部又尝试去调用主题的attach或detach方法即回调中修改观察者列表而主题的通知又用了锁就很容易形成循环等待导致死锁。性能瓶颈如果直接用一把大锁std::mutex保护整个Subject那么任何注册、注销或通知操作都会串行化。在高频通知场景下这锁的争用会成为巨大的性能瓶颈。通知顺序与一致性多线程并发通知下观察者接收到消息的顺序可能和事件发生的真实顺序不一致。某些场景下如状态机切换顺序错乱是灾难性的。生命周期管理观察者对象可能在其他线程被销毁。如果主题不知情仍持有其指针并调用update就会访问已释放的内存导致段错误。因此我们的设计目标很明确在满足线程安全的前提下尽可能保证高性能和可扩展性同时妥善处理对象的生命周期。2.2 总体架构选型几种常见的线程安全观察者模式实现面对上述挑战有几种主流的设计思路各有优劣1. 粗粒度锁Coarse-grained Locking这是最直接的想法在Subject内部用一个std::mutex在attach、detach、notify每个方法的开始都std::lock_guard锁住整个列表操作。优点实现简单线程安全。缺点性能差所有操作串行。特别是在notify时锁会贯穿所有观察者的update调用如果某个观察者处理慢会阻塞所有其他观察者和后续的注册/注销操作。适用场景观察者数量少更新频率极低或者对性能不敏感的原型阶段。2. 写时复制Copy-on-Write核心是维护一个std::shared_ptrconst std::vectorObserver*指向当前的观察者列表。当需要修改列表attach/detach时先复制一份当前列表的副本在副本上修改然后通过原子操作将shared_ptr切换指向新的副本。notify时只需原子地读取这个shared_ptr然后遍历它指向的只读的列表副本。优点notify操作完全无锁读操作性能极高。读写分离清晰。缺点attach/detach的写操作成本较高需要复制整个列表。不适合观察者列表频繁变动的场景。内存开销稍大。适用场景读多通知极多写少注册注销很少的典型场景比如配置变更通知、系统信号广播。3. 细粒度锁与通知队列Fine-grained Locking with Notification Queue这是更复杂但也更强大和灵活的模式。它解耦了“状态变更”和“执行通知”。Subject内部维护一个线程安全的观察者容器如用读写锁std::shared_mutex保护的std::vector或并发容器。当状态变化需要通知时Subject并不直接调用观察者而是将一个“通知任务”包含事件数据和观察者列表的快照或引用放入一个线程安全的队列如moodycamel::ConcurrentQueue或std::queue 锁。由一个或多个专用的工作线程或线程池从队列中取出任务并实际执行对各个观察者的update调用。优点彻底解耦主题线程不会被慢观察者阻塞。可以通过工作线程池控制并发度。灵活支持异步、延迟通知。观察者列表的修改和任务的入队可以设计得更高效。缺点架构复杂引入了消息队列和线程管理。事件传递有延迟队列排队。需要处理任务队列的积压和消费问题。适用场景高性能、高并发、观察者处理耗时不均的工业级系统如金融行情分发、游戏事件系统、微服务间的消息总线。对于大多数需要严肃对待性能的C项目“写时复制”和“通知队列”是更值得考虑的方案。下面我们将深入“写时复制”方案的实现细节因为它平衡了复杂度、安全性和性能是很多开源框架和中间件的基础实现方式。3. 核心实现基于“写时复制”的线程安全观察者3.1 线程安全的观察者列表管理“写时复制”的精髓在于分离读和写。我们使用std::shared_ptr来管理观察者列表利用其原子引用计数的特性来实现无锁的读。首先定义观察者接口和主题基类#include memory #include vector #include mutex #include atomic class Observer { public: virtual ~Observer() default; // 使用const传递消息避免拷贝开销。可根据需要改为模板或any。 virtual void update(const std::string event_data) 0; }; class ThreadSafeSubject { public: virtual ~ThreadSafeSubject() default; virtual void attach(std::weak_ptrObserver obs) 0; virtual void detach(std::weak_ptrObserver obs) 0; // 通常通过token或直接比较来注销 virtual void notify(const std::string event_data) 0; };关键点在于attach和detach接收的是std::weak_ptrObserver。这是解决生命周期问题的核心主题只持有观察者的弱引用避免阻止观察者被销毁。当需要通知时尝试将weak_ptr提升lock()为shared_ptr如果提升成功说明对象还活着可以安全调用。接下来是实现类class CopyOnWriteSubject : public ThreadSafeSubject { // 使用shared_ptr管理一个const的观察者列表。 // mutable 用于在const成员函数(如notify)中修改ptr_这是原子操作。 mutable std::shared_ptrconst std::vectorstd::weak_ptrObserver observers_ptr_; // 写操作attach/detach需要互斥防止同时修改。 mutable std::mutex writers_mutex_; public: CopyOnWriteSubject() : observers_ptr_(std::make_sharedconst std::vectorstd::weak_ptrObserver()) { // 初始化为一个空的常量vector } void attach(std::weak_ptrObserver obs) override { std::lock_guardstd::mutex lock(writers_mutex_); // 1. 复制当前列表解const auto new_observers std::make_sharedstd::vectorstd::weak_ptrObserver(*observers_ptr_); // 2. 修改副本 new_observers-push_back(obs); // 3. 原子地切换指针std::shared_ptr的赋值是线程安全的 observers_ptr_ std::move(new_observers); // 注意这里move后new_observers变为空 } void detach(std::weak_ptrObserver obs) override { // 如何比较weak_ptr通常需要Observer对象提供一个唯一的标识符如ID // 或者我们存储的是weak_ptr无法直接比较指向的对象是否相同。 // 一种常见做法是attach时返回一个token如迭代器或IDdetach时使用这个token。 // 这里展示一个简化版遍历查找并删除。生产环境需要更高效的方案。 std::lock_guardstd::mutex lock(writers_mutex_); auto new_observers std::make_sharedstd::vectorstd::weak_ptrObserver(*observers_ptr_); // 移除失效的或指定的观察者。注意weak_ptr比较需要转为shared_ptr或比较控制块地址。 // 这只是一个示意性删除逻辑。 new_observers-erase( std::remove_if(new_observers-begin(), new_observers-end(), [obs](const std::weak_ptrObserver w) { // 比较控制块地址原始指针是否相同。这要求传入的obs必须是由同一个shared_ptr生成的weak_ptr。 // 更稳健的做法是给Observer一个UUID。 return !(w.owner_before(obs) || obs.owner_before(w)); // owner_before用于比较所有权 }), new_observers-end()); observers_ptr_ std::move(new_observers); } void notify(const std::string event_data) override { // 关键步骤原子地获取当前列表的快照。这个操作本身是线程安全且无锁的。 auto local_copy std::atomic_load(observers_ptr_); // 使用atomic_load显式说明 // 注意local_copy是一个shared_ptr指向一个const vector。后续的遍历是只读的。 for (const auto weak_obs : *local_copy) { // 尝试提升为shared_ptr if (auto obs weak_obs.lock()) { // 提升成功对象存活安全调用 obs-update(event_data); } // 如果提升失败说明观察者对象已被销毁跳过即可。 } } };注意上面的detach实现是示意性的直接比较weak_ptr的 ownership 并不直观且效率低。在实际项目中attach方法通常会返回一个代表注册关系的“令牌”比如一个uint64_t类型的ID或者一个不透明的Registration对象detach时使用这个令牌来精确移除效率更高也避免了复杂的比较逻辑。3.2 通知过程的性能与一致性考量在notify函数中我们通过std::atomic_load或者直接赋值因为shared_ptr的读操作本身是原子的获得了一份观察者列表的“快照”。这份快照在通知过程中是恒定不变的即使其他线程在此期间调用了attach或detach修改的也是observers_ptr_指向的新列表不会影响我们正在遍历的这份旧列表。这保证了遍历过程的安全性。然而这里有几个重要的性能和一致性细节“过期”通知我们遍历的可能是“上一刻”的观察者列表。如果一个观察者在事件发生后、但在我们获取快照前一刻才注册它就会错过这次通知。这在大多数事件系统中是可接受的因为观察者应该关注“未来”的事件。如果需要严格“捕获”所有事件则需要更复杂的同步机制。通知期间的观察者销毁这是weak_ptr解决的经典问题。即使在遍历过程中某个观察者对象在其他线程被销毁了我们的weak_ptr.lock()也会失败从而安全地跳过调用避免了悬空指针。异常安全如果某个观察者的update方法抛出异常会中断整个通知循环。你需要决定是捕获异常并继续通知其他观察者try-catch在循环内还是让异常传播出去中断通知。通常一个观察者的失败不应影响其他观察者。性能热点虽然遍历本身无锁但如果观察者数量巨大成千上万且update调用非常快那么复制整个weak_ptr向量在attach/detach时的成本就会凸显。对于列表频繁变动的场景此方案不适用。一个改进的notify实现增加了异常处理和轻度优化void notify(const std::string event_data) override { auto local_copy std::atomic_load(observers_ptr_); // 快照 // 预分配避免在循环中多次分配如果update是重操作这点优化微乎其微 std::vectorstd::shared_ptrObserver valid_observers; valid_observers.reserve(local_copy-size()); // 第一遍收集所有有效的观察者 for (const auto weak_obs : *local_copy) { if (auto obs weak_obs.lock()) { valid_observers.push_back(std::move(obs)); } } // 第二遍通知有效的观察者并隔离异常 for (const auto obs : valid_observers) { try { obs-update(event_data); } catch (const std::exception e) { // 日志记录错误但继续通知其他观察者 // std::cerr Observer update failed: e.what() std::endl; } catch (...) { // 处理未知异常 // std::cerr Observer update failed with unknown exception. std::endl; } } }4. 高级话题应对复杂场景与性能优化4.1 支持泛型事件与数据传递基础的update(const std::string)接口限制了事件数据的类型。一个更强大的观察者模式应该支持泛型事件。我们可以借助std::any、自定义事件基类或者模板技术。方法一使用std::any(C17)class Observer { public: virtual ~Observer() default; virtual void update(const std::any event_data) 0; };主题的notify可以传递任意类型数据。观察者内部需要用std::any_castT来获取数据。优点是灵活缺点是类型不安全转换错误会导致异常且有一定运行时开销。方法二模板化主题编译时多态templatetypename EventData class Observer { public: virtual ~Observer() default; virtual void update(const EventData event_data) 0; }; templatetypename EventData class CopyOnWriteSubject { // ... 内部容器存储 ObserverEventData* 或 std::weak_ptrObserverEventData // attach, detach, notify 针对 EventData 类型 };这种方法类型安全性能好但会导致代码膨胀并且不同事件类型的观察者无法注册到同一个主题上灵活性降低。方法三类型擦除与自定义事件对象推荐定义一个通用的事件基类然后派生具体事件。这是大型框架如Qt常用的方法。struct Event { virtual ~Event() default; // 可以添加事件类型枚举、时间戳等通用字段 int type_id; }; struct PriceUpdateEvent : public Event { std::string symbol; double price; // ... 其他字段 }; class Observer { public: virtual ~Observer() default; virtual void onEvent(const Event e) 0; };主题通知时传递Event的引用或智能指针。观察者通过type_id或dynamic_cast来判断和处理具体事件。它在灵活性、类型安全和性能之间取得了较好的平衡。4.2 异步通知与线程池集成“写时复制”模式下的notify仍然是同步的即调用notify的线程会阻塞直到所有观察者的update方法执行完毕。对于耗时较长的观察者这可能会阻塞主题线程。我们可以引入一个线程安全的任务队列和线程池将通知异步化#include thread #include queue #include functional #include condition_variable class AsyncNotificationDispatcher { std::queuestd::functionvoid() tasks_; mutable std::mutex queue_mutex_; std::condition_variable cv_; std::vectorstd::thread workers_; bool stop_{false}; public: AsyncNotificationDispatcher(size_t num_threads std::thread::hardware_concurrency()) { for (size_t i 0; i num_threads; i) { workers_.emplace_back([this] { this-workerLoop(); }); } } ~AsyncNotificationDispatcher() { { std::lock_guardstd::mutex lock(queue_mutex_); stop_ true; } cv_.notify_all(); for (auto t : workers_) { if (t.joinable()) t.join(); } } templatetypename F void post(F task) { { std::lock_guardstd::mutex lock(queue_mutex_); tasks_.emplace(std::forwardF(task)); } cv_.notify_one(); } private: void workerLoop() { while (true) { std::functionvoid() task; { std::unique_lockstd::mutex lock(queue_mutex_); cv_.wait(lock, [this] { return stop_ || !tasks_.empty(); }); if (stop_ tasks_.empty()) return; task std::move(tasks_.front()); tasks_.pop(); } task(); // 执行通知任务 } } }; // 在Subject中使用 class AsyncCopyOnWriteSubject : public CopyOnWriteSubject { AsyncNotificationDispatcher dispatcher_; public: AsyncCopyOnWriteSubject(AsyncNotificationDispatcher disp) : dispatcher_(disp) {} void notify(const std::string event_data) override { auto local_copy std::atomic_load(observers_ptr_); // 收集有效的观察者同上 std::vectorstd::shared_ptrObserver valid_observers; for (const auto weak_obs : *local_copy) { if (auto obs weak_obs.lock()) { valid_observers.push_back(std::move(obs)); } } // 将通知任务提交到线程池 dispatcher_.post([valid_observers std::move(valid_observers), event_data]() mutable { for (auto obs : valid_observers) { try { obs-update(event_data); } catch (...) { /* 处理异常 */ } } }); } };这样notify调用将立即返回实际的通知工作由线程池中的工作线程异步执行。注意异步化带来了事件顺序的不确定性如果事件顺序重要需要在事件数据中加入序列号或者使用单线程的消费者模型。4.3 避免常见陷阱死锁、竞态与生命周期在观察者回调中修改主题这是死锁的经典来源。如果观察者的update方法里又调用了subject.attach/detach而主题的内部锁如writers_mutex_还未释放就会导致死锁。黄金法则观察者的回调函数应尽可能快且避免回调中再去同步调用主题的注册/注销方法。如果必须修改应使用异步方式如将修改请求放入队列。weak_ptr失效与性能频繁地对大量weak_ptr调用lock()有一定开销。如果观察者生命周期很长可以考虑存储shared_ptr但这会阻止观察者被自动销毁必须手动detach。通常使用weak_ptr是更安全的选择但需要在设计上接受偶尔的lock开销。“通知风暴”与流量控制如果主题状态变化极快会产生海量通知任务可能压垮观察者或任务队列。需要考虑节流Throttling或采样Sampling。例如可以设置一个最小通知间隔或者只通知最后一次状态。跨线程传递与内存序在我们使用std::atomic_load和shared_ptr赋值时默认的内存序memory_order_seq_cst保证了最强的顺序一致性但可能影响性能。在x86等强内存模型架构上使用memory_order_acquire/memory_order_release可能就足够了。但这属于高级优化需要对C内存模型有深刻理解否则容易引入极难调试的bug。对于大多数应用默认设置是最安全的选择。5. 实战一个简单的股票行情通知系统示例让我们用一个简化版的股票行情通知系统来串联上述概念。系统有一个行情源MarketDataSubject多个策略StrategyObserver订阅特定股票的行情更新。#include iostream #include string #include memory #include unordered_map // 事件数据行情更新 struct QuoteEvent { std::string symbol; double bid; double ask; uint64_t timestamp; }; // 观察者接口 class QuoteObserver { public: virtual ~QuoteObserver() default; virtual void onQuoteUpdate(const QuoteEvent quote) 0; }; // 线程安全的主题基于写时复制简化版 class QuoteSubject { using ObserverPtr std::weak_ptrQuoteObserver; using ObserverList std::vectorObserverPtr; std::shared_ptrconst ObserverList observers_ptr_{std::make_sharedconst ObserverList()}; mutable std::mutex writers_mutex_; public: // 返回一个注册令牌这里用智能指针本身用于后续注销 std::shared_ptrQuoteObserver attach(std::shared_ptrQuoteObserver obs) { std::lock_guardstd::mutex lock(writers_mutex_); auto new_list std::make_sharedObserverList(*observers_ptr_); new_list-push_back(obs); // 存储weak_ptr observers_ptr_ std::move(new_list); return obs; // 调用者保存这个shared_ptr既是观察者对象也是注销令牌 } void detach(const std::shared_ptrQuoteObserver token) { std::lock_guardstd::mutex lock(writers_mutex_); auto new_list std::make_sharedObserverList(); new_list-reserve(observers_ptr_-size()); // 复制除了token对应的weak_ptr之外的所有观察者 for (const auto weak_obs : *observers_ptr_) { if (auto sp weak_obs.lock()) { if (sp ! token) { // 比较shared_ptr找到要删除的 new_list-push_back(weak_obs); } } else { // 已经失效的也一并清理 } } observers_ptr_ std::move(new_list); } void notify(const QuoteEvent quote) { auto local_copy observers_ptr_; // 原子读 for (const auto weak_obs : *local_copy) { if (auto obs weak_obs.lock()) { try { obs-onQuoteUpdate(quote); } catch (const std::exception e) { std::cerr Strategy error on quote.symbol : e.what() std::endl; } } } } }; // 具体的策略观察者 class MovingAverageStrategy : public QuoteObserver, public std::enable_shared_from_thisMovingAverageStrategy { std::string symbol_; double sum_{0.0}; int count_{0}; public: explicit MovingAverageStrategy(std::string sym) : symbol_(std::move(sym)) {} void onQuoteUpdate(const QuoteEvent quote) override { if (quote.symbol symbol_) { sum_ (quote.bid quote.ask) / 2.0; count_; std::cout [MA Strategy] Symbol: symbol_ , Avg Price: (sum_ / count_) (from count_ updates) std::endl; } } }; class AlertStrategy : public QuoteObserver, public std::enable_shared_from_thisAlertStrategy { std::string symbol_; double threshold_; public: AlertStrategy(std::string sym, double th) : symbol_(std::move(sym)), threshold_(th) {} void onQuoteUpdate(const QuoteEvent quote) override { if (quote.symbol symbol_ quote.bid threshold_) { std::cout [Alert] Symbol symbol_ bid price quote.bid exceeds threshold threshold_ ! std::endl; } } }; int main() { QuoteSubject subject; // 创建策略 auto ma_strategy std::make_sharedMovingAverageStrategy(AAPL); auto alert_strategy std::make_sharedAlertStrategy(AAPL, 150.0); // 注册策略并保存返回的令牌这里就是策略对象本身 auto token1 subject.attach(ma_strategy); auto token2 subject.attach(alert_strategy); // 模拟行情更新可以在另一个线程中 subject.notify(QuoteEvent{AAPL, 148.5, 148.6, 1234567890}); subject.notify(QuoteEvent{GOOGL, 2700.1, 2700.3, 1234567891}); // 这个不会被策略处理 subject.notify(QuoteEvent{AAPL, 151.2, 151.3, 1234567892}); // 注销一个策略 subject.detach(token1); subject.notify(QuoteEvent{AAPL, 152.0, 152.1, 1234567893}); // 只有Alert策略会响应 return 0; }这个示例展示了核心流程线程安全的注册注销、基于事件数据的过滤、以及基本的错误处理。在实际系统中行情源notify的调用很可能来自网络I/O线程而策略的计算可能耗时此时就需要引入前面提到的异步通知分发器将notify中的循环转移到线程池中执行。6. 性能测试与选型建议选择哪种实现方式最终要落到性能数据上。我们可以用简单的基准测试来对比测试场景10万个观察者主线程频繁调用notify同时有少量线程随机进行attach/detach。对比方案粗粒度锁std::mutex保护整个向量。写时复制如上文实现。无锁队列如moodycamel::ConcurrentQueue将观察者列表的修改和通知任务完全解耦。预期结果粗粒度锁在纯通知压力下由于锁竞争性能会随着线程数增加而急剧下降。写时复制通知性能极高无锁读且几乎不受读者线程数影响。但写操作注册/注销成本高且随观察者数量线性增长。无锁队列无论是通知还是注册注销吞吐量都可能最高因为它将争用分散到了队列操作上。但延迟会比写时复制略高因为需要任务入队出队且架构复杂度最高。选型建议如果你的场景是观察者数量固定或变化极少如插件系统、配置监听但事件通知极其频繁如实时渲染、游戏循环。选择写时复制。它的无锁读特性是巨大优势。如果你的场景是观察者频繁动态增删如聊天室用户进出且观察者的回调函数执行时间不确定有长有短。选择通知队列线程池。它提供了最好的隔离性和扩展性。如果你的场景是原型验证、观察者很少100、性能要求不高。选择粗粒度锁。简单可靠快速上线。永远记住在实现任何优化前先用性能分析工具如perf、VTune找到真正的热点。过早优化是万恶之源。线程安全观察者模式的复杂性只有在确实面临并发瓶颈时才值得引入。最后多线程环境下的观察者模式其核心矛盾始终是数据共享与同步的开销。没有银弹最好的模式是那个最能贴合你具体业务负载、团队熟悉度和长期维护成本的模式。从简单的加锁开始随着量测和需求演进再逐步升级到更复杂的无锁或队列结构是一个稳健的工程实践路径。