【问题标题】:Two waiting threads (producer/consumer) with a shared buffer具有共享缓冲区的两个等待线程(生产者/消费者)
【发布时间】:2018-09-13 09:31:36
【问题描述】:

我试图让一堆生产者线程等到缓冲区有空间容纳一个项目,然后它会尽可能将项目放入缓冲区,如果没有更多空间则返回睡眠。

同时应该有一堆消费者线程等待直到缓冲区中有东西,然后尽可能从缓冲区中取出东西,如果它是空的则返回睡眠。

在伪代码中,这就是我正在做的事情,但我得到的只是死锁。

condition_variable cvAdd;
condition_variable cvTake;
mutex smtx;

ProducerThread(){
    while(has something to produce){

         unique_lock<mutex> lock(smtx);
         while(buffer is full){
            cvAdd.wait(lock);
         }
         AddStuffToBuffer();
         cvTake.notify_one();
    }
}

ConsumerThread(){

     while(should be taking data){

        unique_lock<mutex> lock(smtx);
        while( buffer is empty ){
            cvTake.wait(lock);
        }   
        TakeStuffFromBuffer();
        if(BufferIsEmpty)
        cvAdd.notify_one();
     }

}

【问题讨论】:

  • 你在哪里遇到死锁?你会遇到什么样的僵局?我的意思是:线程在哪里等待? pthread_cond_wait、pthread_mutex_lock、pthread_join...?我试过你的伪代码,它对我有用,如果你愿意,我可以发布代码。
  • 我在下面完全改变了我原来的答案。

标签: c++ multithreading thread-safety producer-consumer


【解决方案1】:

另一个值得一提的错误是,您的消费者仅在缓冲区变空时才通知等待的生产者。

通知消费者的最佳方式是仅在队列已满时。

例如:

template<class T, size_t MaxQueueSize>
class Queue
{
    std::condition_variable consumer_, producer_;
    std::mutex mutex_;
    using unique_lock = std::unique_lock<std::mutex>;

    std::queue<T> queue_;

public:
    template<class U>
    void push_back(U&& item) {
        unique_lock lock(mutex_);
        while(MaxQueueSize == queue_.size())
            producer_.wait(lock);
        queue_.push(std::forward<U>(item));
        consumer_.notify_one();
    }

    T pop_front() {
        unique_lock lock(mutex_);
        while(queue_.empty())
            consumer_.wait(lock);
        auto full = MaxQueueSize == queue_.size();
        auto item = queue_.front();
        queue_.pop();
        if(full)
            producer_.notify_all();
        return item;
    }
};

【讨论】:

  • 关于您对生产者的通知:可能是缓冲区已满,一些生产者去等待简历。然后安排了一些消费者并每个人弹出一个项目,但只有一个生产者会被唤醒,而其他人仍然在睡觉,不是吗?所以最好在物品被弹出时通知生产者,或者我错过了什么?
  • @MikevanDyke 你是对的,我的逻辑有问题。已更新。
  • @MikevanDyke 工作线程需要存在,直到有一些读者并且队列中有东西。我在一边有一个受互斥锁保护的读卡器计数器,但是如何安全地检查工作线程中的队列是否为空?
  • @tugayicoru 这是一个错误的问题。购买的时候你检查queue.empty()的结果,结果可能不再符合现实。相反,您应该以这样一种方式构建您的代码,即当发布者关闭时,应用程序会等待所有消费者终止。这可以通过不同的方式来完成。一种方法是将一个特殊项目发布到队列中,指示消费者必须终止。
【解决方案2】:

您的生产者和消费者都尝试锁定互斥锁,但没有一个线程解锁互斥锁。这意味着第一个获取锁的线程持有它,而另一个线程永远不会运行。

考虑将互斥锁调用移动到每个线程执行其操作之前,然后在每个线程执行其操作(AddStuffTobuffer() 或 TakeStuffFromBuffer())之后解锁。

【讨论】:

  • 在添加到缓冲区之前,我会检查缓冲区容量。除非我打算同时查看多个线程的缓冲区大小,否则无法真正移动锁。此外,锁在离开作用域时会解锁,因此在执行另一轮 while 循环时会解锁。
  • 您确定循环的每次迭代都是一个新范围吗?我不相信这是真的。范围迭代不会创建新的代码块,这将创建一个新的范围。我看到您的生产者明确调用 lock.unlock() 但您的消费者没有。如果消费者从不释放它的锁,一旦缓冲区清空,程序就会死锁。
  • 我确定。离开循环调用析构函数,它解锁了锁,然后我用另一个迭代再次创建锁。
【解决方案3】:

根据您的查询查看此示例。在这种情况下,一个单独的 condition_variable 就足够了。

#include "conio.h"
#include <thread>
#include <mutex>
#include <queue>
#include <chrono>
#include <iostream>
#include <condition_variable>

using namespace std;

mutex smtx;
condition_variable cvAdd;
bool running ;
queue<int> buffer;

void ProducerThread(){
    static int data = 0;
    while(running){
        unique_lock<mutex> lock(smtx);
        if( !running) return;
        buffer.push(data++);
        lock.unlock();
        cvAdd.notify_one();
        this_thread::sleep_for(chrono::milliseconds(300));
    }
}

void ConsumerThread(){

     while(running){

        unique_lock<mutex> lock(smtx);
        cvAdd.wait(lock,[](){ return !running || !buffer.empty(); });
         if( !running) return;
        while( !buffer.empty() )
        {
            auto data = buffer.front();
            buffer.pop();
            cout << data <<" \n";

            this_thread::sleep_for(chrono::milliseconds(300)); 
        }                

     }

}

int main()
{
    running = true;
    thread producer = thread([](){ ProducerThread(); }); 
    thread consumer = thread([](){ ConsumerThread(); });

    while(!getch())
    { }    

    running = false;
    producer.join();
    consumer.join();  
}

【讨论】:

  • 该程序用于学校作业。我不能使用 chrono 库。
【解决方案4】:

我之前回答过这个问题,但我有点跑题了,因为我目前正在了解mutex、lock_guard 等的基本机制和行为。我一直在观看一些关于我目前正在观看的主题和一个视频实际上与locking 相反,因为视频显示了如何实现LockFreeQueue,它使用循环缓冲区或环形缓冲区、两个指针,并使用atomic 而不是@ 987654331@。现在,对于您目前的情况,atomic 和 LockFreeQueue 将无法回答您的问题,但我从该视频中获得的是循环缓冲区的想法。

因为您的生产者/消费者线程都将共享同一个内存池。如果生产者 - 消费者线程的比例为 1 比 1,则很容易跟踪数组或每个指针的每个索引。但是,当您拥有多对多时,事情确实会变得有些复杂。

可以做的一件事是,如果您将缓冲区的大小限制为 N 个对象,您实际上可能希望将其创建为 N+1。一个额外的空白空间,这将有助于减轻在多个生产者和消费者之间共享的环形缓冲区结构中的一些复杂性。


请看下图:

p = 生产者索引,c = 消费者索引,N 表示 [ ] 索引空间的数量。 N = 5。

一对一

 p                N = 5
[ ][ ][ ][ ][ ]
 c

这里 p 和 c == 0。这表示缓冲区是空的。假设生产者在 c 收到任何内容之前填充缓冲区

             p    N = 5
[x][x][x][x][x]
 c

在这种情况下,缓冲区已满,p 必须等待一个空白空间。 c 现在可以获取了。

             p     N = 5
[ ][x][x][x][x]
    c         

这里 c 在 [0] 处获取对象,并将其索引提升到 1。 P 现在可以在环形缓冲区周围转一圈了。

这很容易通过单个 p & c 来跟踪。现在让我们探索多个消费者和单个生产者

一对多

 p                 N = 5
[ ][ ][ ][ ][ ]
c1
c2

这里p index = 0,c1 & c2 index = 0,环形缓冲区是空的

             p     N = 5
[x][x][x][x][x]
c1
c2

现在 p 必须等待 c1 或 c2 在 [0] 处获取项目才能写入

             p     N = 5
[ ][ ][x][x][x]
    c1 c2

这里不清楚 c1 或 c2 是否获得了 [0] 或 1,但两者都成功获得了一个项目。两者都增加了索引计数器。以上似乎表明 c1 从 [0] 增加到 1。然后在 [0] 处的 c2 也必须增加索引计数器,但它已经从 0 更改为 1,因此 c2 将其增加为 2。

如果我们假设当p == 0 &amp;&amp; c1 || c2 == 0 缓冲区为空时,这里就会出现死锁情况。看看这里的情况。

 p               N = 5  // P hasn't written yet but has advanced 
[ ][ ][ ][ ][x]  // 1 Item is left
           c1  // Both c1 & c2 have index to same item.
           c2  // c1 acquires it and so does c2 but one of them finishes first and then increments the counter. Now the buffer is empty and looks like this:

 p                N = 5
[ ][ ][ ][ ][ ]
c1          c2    // p index = 0 and c1 = 0 represents empty buffer.
                  // c2 is trying to read [4]

这会导致死锁。

多对一

 p1
 p2                N = 5
[ ][ ][ ][ ][ ]
 c1

这里有多个生产者可以为单个消费者写入缓冲区。如果它们交错:

p1 writes to [0] increments counter
p2 writes to [0] increments counter

   p1 p2
[x][ ][ ][ ][ ]
c1

这将导致缓冲区中出现空白空间。生产者互相干扰。这里需要互斥。

多对多的想法;您需要考虑并结合上述一对多和多对一的两个功能。您需要为消费者使用互斥锁,为生产者使用互斥锁,尝试为两者使用相同的互斥锁会给您带来可能导致无法预料的死锁的问题。您必须确保检查所有案件并在适当的时间锁定它们 - 地点。也许这几个视频可以帮助你了解更多。


伪代码:可能如下所示:

condition_variable cvAdd;
condition_variable cvTake;
mutex consumerMutex;
mutex producerMutex;

ProducerThread(){
    while( has something to produce ) {    
         unique_lock<mutex> lock(producerMutex);
         while(buffer is full){
            cvAdd.wait(lock);
         }
         AddStuffToBuffer();
         cvTake.notify_one();
    }
}

ConsumerThread() {    
     while( should be taking data ) {    
        unique_lock<mutex> lock(consumerMutex);
        while( buffer is empty ){
            cvTake.wait(lock);
        }   
        TakeStuffFromBuffer();
        if(BufferIsEmpty)
        cvAdd.notify_one();
     }    
}

这里唯一的区别是使用了 2 个独占互斥体,而不是生产者和消费者都尝试使用同一个互斥体。共享的是内存;但是您不想在两者之间共享计数器或指针到内存池中。多个生产者使用同一个互斥锁是可以的,多个消费者使用同一个互斥锁也可以,但是让消费者和生产者都使用同一个互斥锁可能是您的根本问题。

【讨论】:

  • 不说为什么就投了反对票?多么典型...留下一个简洁的理由说明为什么反对票让位于能够进行编辑以更新原始答案以提高其准确性。
  • 在您的示例中,std::shared_lock 看起来并没有保护任何东西,因为它允许多个线程访问锁定下的数据。这没关系,只要所有线程只读取信息但不更改它。如果您想更改信息,您必须创建一个std::unique_lock。
  • @Galik 好的;谢谢你澄清这一点。在我自己的一个类中,我尝试使用lock_guard&lt;mutex&gt;,与上面显示的相同,代码编译和构建,但是当我尝试运行它时,它会引发异常。我也以同样的方式尝试了unique_lock&lt;mutex&gt;,做了同样的事情,它编译和构建,但它抛出了一个异常。当我将它们更改为 shared_lock&lt;shared_mutex&gt; 时,据我所知,代码的行为正常。
  • 该程序用于学校作业。我不能使用 shared_mutex 库。
  • @tugayicoru 实际上;我最近开始使用互斥体,并且开始更好地理解它们,shared_mutex 实际上是一个错误的实现。
猜你喜欢
  • 2011-02-15
  • 1970-01-01
  • 1970-01-01
  • 2019-01-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-02-16
  • 2016-06-29
相关资源
最近更新 更多