如何移除多线程发布/订阅实现中的recursive_mutex?
如何在多线程发布/订阅实现中移除recursive_mutex?
我知道使用recursive_mutex被视为代码异味,通常存在非递归互斥锁的等效实现。现在需要在这个多线程发布/订阅(Publish/Subscribe)实现里移除recursive_mutex。
PubSub类会根据ContainerProperty(已简化为int类型)将Container对象推送给已订阅的Observer。当前的核心问题是允许Observer在Update()回调中取消订阅:添加UpdateObserver结构体及相关函数是为了避免遍历Observer时执行订阅/取消订阅操作,但仍需要在Unsubscribe()时立即使Observer失效,这就必须对valid属性(ObserverList.second)加锁,防止在检查有效性和调用Update()之间Observer被销毁。
我尝试过用weak_ptr<Observer>替代Observer*的方案,这样就无需在Unsubscribe()时使Observer失效(可以在使用前通过lock()检查对象是否存在)。该方案看似可行,但重构幅度较大,还存在额外问题:我的大多数Observer会在构造函数中订阅,这与enable_shared_from_this()的配合并不顺畅,但这可能是另一个独立问题。
我希望找到直接解决recursive_mutex问题的方法,同时尽可能少地修改现有接口。
相关代码
observer.h
#ifndef _OBSERVER_H_INCLUDED_ #define _OBSERVER_H_INCLUDED_ #include <vector> class PubSub; typedef int ContainerProperty; struct Container { std::vector<double> data; ContainerProperty property; }; class Observer { public: Observer(PubSub& p) : m_pubSub(p) {} virtual ~Observer(); virtual void Update(const Container&) = 0; protected: PubSub& m_pubSub; }; #endif
observer.cpp
#include "observer.h" #include "pubsub.h" Observer::~Observer() { m_pubSub.Unsubscribe(*this); }
pubsub.h
#ifndef _PUBSUB_H_INCLUDED_ #define _PUBSUB_H_INCLUDED_ #include <condition_variable> #include <list> #include <map> #include <mutex> #include <queue> #include <thread> #include "observer.h" //------------------------------------------ // //------------------------------------------ class PubSubInterface { public: virtual ~PubSubInterface() = default; virtual void Publish(Container&& c) = 0; virtual void Subscribe(ContainerProperty p, Observer& o) = 0; virtual void Unsubscribe(Observer& o) = 0; }; //------------------------------------------ // //------------------------------------------ class PubSub : public PubSubInterface { std::condition_variable m_condition; std::queue<Container> m_queue; std::mutex m_queueMutex; std::thread m_thread; bool m_running; public: PubSub() : m_running(false) {} ~PubSub() { StopThread(); std::lock_guard<std::mutex> lock(m_queueMutex); std::queue<Container>().swap(m_queue); } bool StartThread() { if (!m_running) { m_running = true; m_thread = std::thread([=] { PubSubThread(); }); if (!m_thread.joinable()) { perror("Unable to start PubSub thread"); m_running = false; } } return m_running; } void StopThread() { if (m_running) { m_running = false; // Wake the thread to allow it to join Publish(Container()); if (m_thread.joinable()) { m_thread.join(); } } } void Publish(Container&& c) override { { std::lock_guard<std::mutex> lock(m_queueMutex); m_queue.push(std::move(c)); } m_condition.notify_all(); } void Subscribe(ContainerProperty p, Observer& o) override { std::lock_guard<std::mutex> lock(m_subscriptionMutex); m_observerUpdates.push_back(UpdateObserver(p, o)); } void Unsubscribe(Observer& o) override { { std::lock_guard<std::mutex> lock(m_subscriptionMutex); m_observerUpdates.push_back(UpdateObserver(o)); } // Invalidate observer immediately // Lock prevents an observer being invalidated while iterating in PubSub thread std::lock_guard<std::recursive_mutex> lock(m_observerMutex); for (auto& property : m_observers) { for (auto& itr : property.second) { if (itr.first == &o) itr.second = false; } } } void WaitQueueEmpty() { std::unique_lock<std::mutex> lock(m_queueMutex); m_condition.wait(lock, [&]() { return m_queue.empty(); }); } private: Container& front() { std::unique_lock<std::mutex> lock(m_queueMutex); m_condition.wait(lock, [&]() { return !m_queue.empty(); }); return m_queue.front(); } void PubSubThread() { while (m_running) { Container& c = front(); { std::lock_guard<std::recursive_mutex> lock(m_observerMutex); UpdateObservers(); const auto& it = m_observers.find(c.property); if (it != m_observers.cend()) { for (const auto& observer : it->second) { if (observer.second) { // recursive_mutex prevents observer being invalidated // after checking validity and avoids deadlock // if Update() in turn calls Unsubscribe() observer.first->Update(c); } } } } { std::lock_guard<std::mutex> lock(m_queueMutex); m_queue.pop(); } if (m_queue.empty()) m_condition.notify_all(); } } void _subscribe(ContainerProperty p, Observer* o) { m_observers[p][o] = true; } void _unsubscribe(Observer* o) { for (auto& observers : m_observers) { for (auto it = observers.second.begin(); it != observers.second.end();) { if (!it->second && it->first == o) it = observers.second.erase(it); else ++it; } } } void UpdateObservers() { std::lock_guard<std::mutex> lock(m_subscriptionMutex); for (const auto& update : m_observerUpdates) { if (update.m_remove) { _unsubscribe(update.m_observer); } else { _subscribe(update.m_property, update.m_observer); } } m_observerUpdates.clear(); } typedef std::map<Observer*, bool> ObserverMap; std::map<ContainerProperty, ObserverMap> m_observers; std::recursive_mutex m_observerMutex; // Locks access to the ObserverMap struct UpdateObserver { UpdateObserver(ContainerProperty p, Observer& o) : m_property(p), m_observer(&o), m_remove(false) {} UpdateObserver(Observer& o) : m_property(0), m_observer(&o), m_remove(true) {} ContainerProperty m_property; Observer* m_observer; bool m_remove = false; }; std::list<UpdateObserver> m_observerUpdates; std::mutex m_subscriptionMutex; // Locks access to the UpdateObserver list }; #endif // _PUBSUB_H_INCLUDED_
main.cpp
#include "pubsub.h" #include <iostream> class ConcreteObserver : public Observer { public: ConcreteObserver(PubSub& p) : Observer(p) {} private: void Update(const Container& c) override { // Do something with the container // Perhaps unsubscribe or subscribe to other properties std::cout << "Handled Container"; } }; int main() { PubSub impl; impl.StartThread(); ConcreteObserver o(impl); impl.Subscribe(3, o); Container c { {1.0, 1.5, 1.2}, 3 }; impl.Publish(std::move(c)); // Added for convenience // In general I don't care if something is left in the queue at destruction impl.WaitQueueEmpty(); }
编译命令
g++ main.cpp observer.cpp -lpthread
内容的提问来源于stack exchange,提问作者S Meredith
相关产品推荐
相关产品推荐

