You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何移除多线程发布/订阅实现中的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 01:22:01