我之前回答过这个问题,但我有点跑题了,因为我目前正在了解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 && 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 个独占互斥体,而不是生产者和消费者都尝试使用同一个互斥体。共享的是内存;但是您不想在两者之间共享计数器或指针到内存池中。多个生产者使用同一个互斥锁是可以的,多个消费者使用同一个互斥锁也可以,但是让消费者和生产者都使用同一个互斥锁可能是您的根本问题。