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

基于Condition Variable的多生产者多消费者问题:实现与代码修正问询

多线程Web服务器并发问题解答

1. 请求存储方式选择

在固定大小队列、代码极简且不使用异步的前提下,选择选项1(复用槽位的环形队列实现),原因如下:

  • 选项2的"追加至队尾"是线性队列思路,当队列满且有元素被消费后,前面的空槽位无法复用,要么需要移动元素(增加代码复杂度),要么会快速耗尽队列空间,不符合固定大小队列的设计要求。
  • 选项1的复用槽位本质是环形队列,通过维护读写指针并取模,能充分利用固定大小的队列空间,代码实现简洁,不需要额外的元素移动操作,完全适配生产者-消费者模型的固定队列场景。

环形队列逻辑示例:

Queue size = 5
读写指针:write=0, read=0, count=0

Client1放入请求:write=1, count=1 → 队列[R1, _, _, _, _]
Client2放入请求:write=2, count=2 → 队列[R1, R2, _, _, _]
Worker取出请求:read=1, count=1 → 队列[_, R2, _, _, _]
Client1放入新请求:write=0(5取模), count=2 → 队列[R3, R2, _, _, _]

2. Mutex的使用场景

需要使用Mutex的场景

  • 所有访问或修改共享队列资源的操作:包括读写队列元素、修改队列的读写指针/计数变量(如front、rear、count)。
  • 检查队列空/满状态的操作:这些状态依赖共享变量,必须在互斥锁保护下判断,避免多线程竞争导致的状态不一致。
  • 调用pthread_cond_wait之前必须持有Mutex:因为pthread_cond_wait会自动释放锁并进入等待,被唤醒后会重新获取锁,确保等待过程中队列状态不会被其他线程修改。

不需要使用Mutex的场景

  • 线程内部局部变量的操作:比如clientId、workerId这类仅当前线程可见的变量。
  • 不涉及共享资源的代码段:比如模拟请求生成/处理的sleep调用、线程内部的日志打印(不涉及共享变量的打印)。

3. 代码修正及说明

原代码核心问题

  1. 队列实现逻辑错误:消费者线程中修改rear--完全破坏了队列的状态,rear应仅由生产者维护,消费者只需移动front指针。
  2. 队列空满判断逻辑错误:当前实现无法复用已消费的槽位,导致固定大小队列的空间被浪费,且满队列判断条件rear == MAX_REQUESTS-1在有消费操作后永远无法满足,队列实际不会"满"。
  3. 信号量使用冗余:pthread_cond_broadcast会唤醒所有等待线程,容易引发惊群效应,改用pthread_cond_signal唤醒单个线程更高效。

修正后的代码(环形队列实现)

#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>
#include <pthread.h>
#include <time.h>

#define MAX_REQUESTS 10

// 共享资源:环形队列
int requestQueue[MAX_REQUESTS];
int writePtr = 0;   // 生产者写入指针
int readPtr = 0;    // 消费者读取指针
int requestCount = 0; // 当前队列中的请求数

pthread_mutex_t queueMutex = PTHREAD_MUTEX_INITIALIZER;
pthread_cond_t queueNotEmpty = PTHREAD_COND_INITIALIZER;
pthread_cond_t queueNotFull = PTHREAD_COND_INITIALIZER;

// 生产者(客户端)线程
void* client(void* arg) {
    int clientId = *((int*)arg);

    while (1) {
        // 模拟请求生成延迟
        sleep(rand() % 2);
        printf("客户端 %d 生成请求...\n", clientId);

        pthread_mutex_lock(&queueMutex);
        // 等待队列非满
        while (requestCount == MAX_REQUESTS) {
            printf("客户端 %d 等待队列空闲...\n", clientId);
            pthread_cond_wait(&queueNotFull, &queueMutex);
        }

        // 写入请求到队列
        requestQueue[writePtr] = clientId;
        writePtr = (writePtr + 1) % MAX_REQUESTS;
        requestCount++;
        printf("客户端 %d 提交请求,当前队列请求数:%d\n", clientId, requestCount);

        // 唤醒一个等待的消费者线程
        pthread_cond_signal(&queueNotEmpty);
        pthread_mutex_unlock(&queueMutex);
    }
    return NULL;
}

// 消费者(工作线程)线程
void* worker(void* arg) {
    int workerId = *((int*)arg);

    while (1) {
        printf("工作线程 %d 就绪...\n", workerId);

        pthread_mutex_lock(&queueMutex);
        // 等待队列非空
        while (requestCount == 0) {
            printf("工作线程 %d 等待请求...\n", workerId);
            pthread_cond_wait(&queueNotEmpty, &queueMutex);
        }

        // 读取并处理请求
        int clientId = requestQueue[readPtr];
        readPtr = (readPtr + 1) % MAX_REQUESTS;
        requestCount--;
        printf("工作线程 %d 处理客户端 %d 的请求,当前队列请求数:%d\n", workerId, clientId, requestCount);

        // 唤醒一个等待的生产者线程
        pthread_cond_signal(&queueNotFull);
        pthread_mutex_unlock(&queueMutex);

        // 模拟请求处理延迟
        sleep(rand() % 3);
    }
    return NULL;
}

int main() {
    srand(time(NULL));
    pthread_t clients[2];
    pthread_t workers[2];
    int clientIds[2] = {1, 2};
    int workerIds[2] = {1, 2};

    // 创建线程
    for (int i = 0; i < 2; ++i) {
        pthread_create(&clients[i], NULL, client, &clientIds[i]);
        pthread_create(&workers[i], NULL, worker, &workerIds[i]);
    }

    // 等待线程结束(实际会无限运行)
    for (int i = 0; i < 2; ++i) {
        pthread_join(clients[i], NULL);
        pthread_join(workers[i], NULL);
    }

    return 0;
}

修正说明

  • 改用环形队列实现,通过writePtr、readPtr和requestCount维护队列状态,充分复用槽位,符合选项1的设计思路。
  • 修复了队列空满判断逻辑:用requestCount直接判断,避免指针操作的竞态问题。
  • 将pthread_cond_broadcast改为pthread_cond_signal,减少不必要的线程唤醒,避免惊群效应。
  • 移除了消费者线程中错误的rear--操作,队列状态仅由生产者和消费者各自维护自己的指针,逻辑清晰。

内容的提问来源于stack exchange,提问作者Intangible _pg18

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 23:27:23