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

C语言pthreads实现生产者/消费者模式:pthread_mutex行为异常

生产者/消费者模型同步问题的修复方案

你的核心问题是条件变量的同步逻辑没有正确绑定生产者和消费者的状态流转,导致生产者无阻塞地发送所有缓冲区,消费者来不及处理,同时mutex的抢占顺序不符合预期。下面是具体的修复思路和代码示例:

共享状态的正确设计

首先需要明确共享结构体必须包含双向同步的状态标记和对应的条件变量,确保生产者只在消费者准备好时发送下一块数据,消费者只在有新数据时处理:

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

#define BUFFER_COUNT 5
#define BUFFER_SIZE 10000

typedef struct {
    int *current_buf;          // 当前待处理的缓冲区指针
    pthread_mutex_t mutex;
    pthread_cond_t cond_buffer_ready;  // 通知消费者有新缓冲区
    pthread_cond_t cond_consumer_ready;// 通知生产者消费者已就绪
    _Bool buffer_ready;        // 是否有未处理的新缓冲区
    _Bool consumer_ready;      // 消费者是否准备好接收下一个
    int done;                  // 是否所有数据已发送完毕
} SharedData;

生产者线程实现

生产者需要严格等待消费者就绪后再发送数据,并且每发送一块就等待消费者处理完成的信号:

void *producer_thread(void *arg) {
    SharedData *data = (SharedData *)arg;

    // 先等待消费者线程启动并就绪
    pthread_mutex_lock(&data->mutex);
    while (!data->consumer_ready) {
        pthread_cond_wait(&data->cond_consumer_ready, &data->mutex);
    }
    pthread_mutex_unlock(&data->mutex);

    for (int i = 0; i < BUFFER_COUNT; i++) {
        // 生成缓冲区数据:填充 (i*1000 + 1) 到 (i*1000 + 10000)
        int *buf = malloc(BUFFER_SIZE * sizeof(int));
        for (int j = 0; j < BUFFER_SIZE; j++) {
            buf[j] = i * 1000 + j + 1;
        }

        // 等待消费者准备好接收下一块数据
        pthread_mutex_lock(&data->mutex);
        while (!data->consumer_ready) {
            pthread_cond_wait(&data->cond_consumer_ready, &data->mutex);
        }
        // 更新缓冲区并标记为就绪
        data->current_buf = buf;
        data->buffer_ready = 1;
        data->consumer_ready = 0; // 消费者开始处理,暂时不再就绪
        pthread_cond_signal(&data->cond_buffer_ready);
        pthread_mutex_unlock(&data->mutex);
    }

    // 通知消费者所有数据已发送完毕
    pthread_mutex_lock(&data->mutex);
    data->done = 1;
    pthread_cond_signal(&data->cond_buffer_ready);
    pthread_mutex_unlock(&data->mutex);

    return NULL;
}

消费者线程实现

消费者需要先标记自己就绪,然后循环等待新缓冲区,处理完成后再次标记就绪并唤醒生产者:

void *consumer_thread(void *arg) {
    SharedData *data = (SharedData *)arg;

    pthread_mutex_lock(&data->mutex);
    data->consumer_ready = 1; // 标记自己已就绪
    pthread_cond_signal(&data->cond_consumer_ready);
    pthread_mutex_unlock(&data->mutex);

    while (1) {
        pthread_mutex_lock(&data->mutex);
        // 等待新缓冲区或任务结束
        while (!data->buffer_ready && !data->done) {
            pthread_cond_wait(&data->cond_buffer_ready, &data->mutex);
        }

        if (data->done) {
            pthread_mutex_unlock(&data->mutex);
            break;
        }

        // 处理缓冲区:读取第100位(索引99)的值
        printf("%d\n", data->current_buf[99]);
        free(data->current_buf); // 释放缓冲区内存

        // 标记处理完成,准备接收下一块
        data->buffer_ready = 0;
        data->consumer_ready = 1;
        pthread_cond_signal(&data->cond_consumer_ready);
        pthread_mutex_unlock(&data->mutex);
    }

    return NULL;
}

主函数初始化与线程启动

int main() {
    SharedData data = {0};
    pthread_mutex_init(&data.mutex, NULL);
    pthread_cond_init(&data.cond_buffer_ready, NULL);
    pthread_cond_init(&data.cond_consumer_ready, NULL);

    pthread_t producer, consumer;
    pthread_create(&consumer, NULL, consumer_thread, &data);
    pthread_create(&producer, NULL, producer_thread, &data);

    pthread_join(producer, NULL);
    pthread_join(consumer, NULL);

    pthread_mutex_destroy(&data.mutex);
    pthread_cond_destroy(&data.cond_buffer_ready);
    pthread_cond_destroy(&data.cond_consumer_ready);

    return 0;
}

关键修复点说明

  1. 双向状态标记:通过buffer_ready和consumer_ready两个标记,严格控制生产者和消费者的流转:生产者必须等消费者就绪才能发数据,消费者必须等有新数据才处理。
  2. 条件变量的正确等待:所有pthread_cond_wait都放在while循环中,避免虚假唤醒(这是POSIX条件变量的标准用法)。
  3. 互斥量的正确范围:仅在修改共享状态或等待条件时持有mutex,避免长时间占用导致线程阻塞。

运行修复后的代码,会稳定输出预期的:

100
1100
2100
3100
4100

内容的提问来源于stack exchange,提问作者robert atwood

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 20:03:11