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

单生产者多消费者无界缓冲区:条件变量实现问题求助

生产者-消费者问题诊断与解决方案

核心问题定位

程序行为不稳定的根源是完全缺失线程同步机制:

  • 多个reader线程同时操作缓冲区head指针,引发数据竞争,导致重复读取、指针越界或死锁
  • reader线程未处理缓冲区为空时的阻塞逻辑,出现空读、无限轮询甚至非法内存访问
  • writer线程操作tail时未与reader同步,当缓冲区为空插入首个节点时需同时修改head和tail,并发场景下会直接破坏链表结构

关于writer仅操作tail的疑问

这个逻辑不完全正确:单链表实现的FIFO缓冲区中,reader线程删除最后一个节点时,也需要修改tail指针(将其置为NULL)。因此writer和reader都会操作tail,必须通过互斥锁同步所有对缓冲区共享资源(head、tail、节点链表)的访问。

具体解决方案

1. 为缓冲区添加同步原语

在缓冲区结构体中加入互斥锁和条件变量,保护共享资源:

#include <pthread.h>

typedef struct sbuffer_node {
    sensor_data_t data;
    struct sbuffer_node *next;
} sbuffer_node_t;

typedef struct sbuffer {
    sbuffer_node_t *head;
    sbuffer_node_t *tail;
    pthread_mutex_t mutex;       // 保护缓冲区所有共享状态
    pthread_cond_t not_empty;    // 缓冲区非空时唤醒等待的reader
    int is_closed;               // 标记writer是否已结束写入
} sbuffer_t;

2. 线程安全的sbuffer_insert实现

int sbuffer_insert(sbuffer_t *buffer, sensor_data_t *data) {
    if (!buffer || !data) return -1;

    sbuffer_node_t *new_node = malloc(sizeof(sbuffer_node_t));
    if (!new_node) return -1;
    new_node->data = *data;
    new_node->next = NULL;

    // 加锁后修改缓冲区状态
    pthread_mutex_lock(&buffer->mutex);
    if (buffer->tail == NULL) {  // 缓冲区为空时同时更新head和tail
        buffer->head = new_node;
        buffer->tail = new_node;
    } else {
        buffer->tail->next = new_node;
        buffer->tail = new_node;
    }
    pthread_cond_signal(&buffer->not_empty);  // 唤醒等待的reader
    pthread_mutex_unlock(&buffer->mutex);

    return 0;
}

3. 线程安全的sbuffer_remove实现

处理空缓冲区阻塞和流结束逻辑:

int sbuffer_remove(sbuffer_t *buffer, sensor_data_t *data) {
    if (!buffer || !data) return -1;

    pthread_mutex_lock(&buffer->mutex);
    // 缓冲区为空且writer未结束,阻塞等待
    while (buffer->head == NULL && !buffer->is_closed) {
        pthread_cond_wait(&buffer->not_empty, &buffer->mutex);
    }

    // writer已结束且缓冲区为空,返回流结束信号
    if (buffer->head == NULL && buffer->is_closed) {
        pthread_mutex_unlock(&buffer->mutex);
        return -2;
    }

    // 取出节点并更新缓冲区状态
    sbuffer_node_t *temp = buffer->head;
    *data = temp->data;
    buffer->head = temp->next;
    if (buffer->head == NULL) {
        buffer->tail = NULL;  // 删除最后一个节点时更新tail
    }

    free(temp);
    pthread_mutex_unlock(&buffer->mutex);
    return 0;
}

4. 修改writer线程逻辑

写入完成后标记结束并唤醒所有reader:

void *writer_thread(void *arg) {
    sbuffer_t *buffer = (sbuffer_t *)arg;
    sensor_data_t data;
    FILE *fp = fopen("sensor_data.txt", "r");

    // 读取文件数据插入缓冲区
    while (read_data_from_file(fp, &data)) {  // 替换为你的文件读取逻辑
        sbuffer_insert(buffer, &data);
    }

    // 写入流结束标识
    data.id = 0;
    sbuffer_insert(buffer, &data);

    // 标记writer结束,唤醒所有等待的reader
    pthread_mutex_lock(&buffer->mutex);
    buffer->is_closed = 1;
    pthread_cond_broadcast(&buffer->not_empty);
    pthread_mutex_unlock(&buffer->mutex);

    fclose(fp);
    pthread_exit(NULL);
}

5. 修改reader线程逻辑

处理流结束并正常退出:

void *reader_thread(void *arg) {
    sbuffer_t *buffer = (sbuffer_t *)arg;
    sensor_data_t data;
    int ret;

    while (1) {
        ret = sbuffer_remove(buffer, &data);
        if (ret == -2) {  // 流结束,退出线程
            break;
        }
        if (ret == 0) {
            printf("Reader %lu: id=%d, value=%.2f\n", pthread_self(), data.id, data.value);
            // 收到流结束标识也可直接退出
            if (data.id == 0) break;
        }
    }

    pthread_exit(NULL);
}

6. 缓冲区初始化与销毁

确保同步原语正确初始化和销毁:

sbuffer_t *sbuffer_init() {
    sbuffer_t *buffer = malloc(sizeof(sbuffer_t));
    if (!buffer) return NULL;
    buffer->head = NULL;
    buffer->tail = NULL;
    buffer->is_closed = 0;
    pthread_mutex_init(&buffer->mutex, NULL);
    pthread_cond_init(&buffer->not_empty, NULL);
    return buffer;
}

void sbuffer_destroy(sbuffer_t *buffer) {
    if (!buffer) return;
    pthread_mutex_lock(&buffer->mutex);
    sbuffer_node_t *temp;
    while (buffer->head) {
        temp = buffer->head;
        buffer->head = temp->next;
        free(temp);
    }
    pthread_mutex_unlock(&buffer->mutex);
    pthread_mutex_destroy(&buffer->mutex);
    pthread_cond_destroy(&buffer->not_empty);
    free(buffer);
}

关键注意事项

  • 所有对head、tail、is_closed的访问必须在互斥锁保护下进行,杜绝数据竞争
  • 使用pthread_cond_wait替代轮询,减少CPU占用并保证线程安全
  • writer结束时必须用pthread_cond_broadcast唤醒所有等待的reader,避免线程永久阻塞
  • 流结束标识配合is_closed标记,确保所有reader都能正确识别结束信号并退出

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 17:05:20