【问题标题】:One producer thread, several consumers一个生产者线程,多个消费者
【发布时间】:2023-03-23 02:55:01
【问题描述】:

我有一个生产者线程和几个消费者,每个消费者都有:拥有数据和唯一 ID 的自己的队列。 我使用 std::map 来识别线程的每个队列。

typedef std::map<int, std::queue<Task>> TaskMap;
TaskMap inputQueue;
TaskMap outputQueue;

每个消费者线程都在处理他的队列中的数据,如果队列为空,线程必须等待数据。 如果我只想用一个线程来做,我可以将 std::condition_variable 与 std::unique_lock 一起使用,但我有几个消费者,所以我需要几个 std::condition_variable,但我不能将它们保存在容器中(复制/赋值是已删除)。 所以我使用这样的代码

while(q.empty()) {
    std::cout << "waiting...\n";
    std::this_thread::sleep_for(std::chrono::milliseconds(100));
}

其中 q 是对队列的引用。 但是我怎样才能用更好的方式同步呢? 提前致谢。 附:队列总会有数据,最后一个数据必须说'exit'。

【问题讨论】:

  • 缺少复制/分配并不意味着您不能将东西放入容器中。只有部分操作不可用。您仍然可以使用emplaceat
  • 请看this
  • “但我有几个消费者,所以我需要几个 std::condition_variable” 不正确。
  • 那么如何识别需要唤醒的线程呢?
  • @user3365834 你在这里真的需要一个线程池。不要试图重新发明轮子。

标签: c++ multithreading c++11


【解决方案1】:

因为每个消费者都有自己的队列,并且所有消费者只有一个生产者,所以本质上是一个生产者一消费者的场景。

换句话说,您没有在所有消费者之间共享一个队列。

【讨论】:

  • 生产者为所有消费者制作数据。在生产者线程中,我有数据和消费者 id,所以我把这些数据放在用 id 标识的队列中。
  • @user3365834 不过,只要消费者不使用 same 队列从中获取数据,它本质上就是单一生产者-单一消费者。有多个单一消费者并且您的生产者为多个消费者扮演单一生产者这一事实并不会改变线程以单一生产者-单一消费者的方式相互同步的事实。
  • @Maxim Egorushkin 据我了解一个线程填充队列并且许多线程可以消耗一个队列的问题
【解决方案2】:

一个std::conditional_variable 和一个std::mutex 就足够了。

Task t;
{
  std::unique_lock<std::mutex> lock(mtx);
  while (q.empty())
    cond_var.wait(lock);
  t = std::move(q.front());
  q.pop_front();
}

主线程会做

{
  std::lock_guard<std::mutex> lock(mtx);
  q.emplace_front(/*...*/);
  cond_var.notify_all();
}

主线程将唤醒所有线程,但大多数线程会重新进入睡眠状态,因为它们的队列仍然是空的。

【讨论】:

  • 但是性能呢?唤醒多线程是 heave 操作,不是吗?
  • 所以使用notify_one() 并且只唤醒一个消费者,它将接受任务并继续工作。
  • 您不能将std::mutex 传递给condition_variable::wait(),并且此示例不是异常安全的。您不应该使用mtx.lock()mtx.unlock(),而是使用unique_lock
  • @JonathanWakely 谢谢,没注意到。
  • 不,你让事情变得更糟了。不要在线程之间共享unique_lock,不要显式调用lock()unlock()。我已经为你修好了。
【解决方案3】:

为此任务调用 OOP:

class ConcurentTaskQueue{
    std::mutex lock;
    std::condition_variable m_ConditionVariable;
    std::queue<Task> m_TaskQueue;


public:
    Task getTask(){
        std::unique_lock<std::mutex> synchLock (lock); //NOTE: consider doing this with a while(programIsRunning){} loop
        while(m_TaskQueue.empty()){
            m_ConditionVariable.wait(synchLock);
        } 
        Task task(std::move(m_TaskQueue.front()));
        m_TaskQueue.pop();
        return task;
    }

    void addTask (Task task){
        std::unique_lock<std::mutex> synchLock (lock);
        m_TaskQueue.push(std::move(task));
        m_ConditionVariable.notify_one();
    }
};

现在简单地说:

std::map<size_t,ConcurentTaskQueue> inputQueue;
std::thread producer ([&]{
      Task task = produceTask()
      inputQueue[ID].addTask(task);
});

std::thread consumer1([&]{
      Task task = inputQueue[ID].getTask();
});


std::thread consumer2([&]{
      Task task = inputQueue[ID].getTask();
});

编辑2: 使用线程池

【讨论】:

  • 您的实现没有考虑wait 上的虚假唤醒,这可能会导致未定义的行为。
  • 你是对的,但这只是方向,不是完整的实现
  • 那么请在答案中包含您的假设。损坏的代码比没有代码更糟糕。
  • 您不能通过添加原子来解决虚假唤醒,而是通过在谓词上循环来解决虚假唤醒。在这种情况下,它应该是while (m_TaskQueue.empty()) 而不是if (m_TaskQueue.empty())
  • return std::move(task); 错误,它阻止了 RVO,return task; 仍然可以使用移动构造函数但也允许 RVO。
【解决方案4】:

如果您这样做不是为了研究/学习,而是为了生产代码,我建议您使用众多实现中的一种。正确、高效和可扩展地实现这一点并不容易,但其他人已经解决了这些问题。有很多选择,例如

【讨论】:

  • 鉴于此处的答案中建议的可怕的损坏代码,我倾向于同意。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-02-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多