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

多写多读模式下Readers-Writers问题实现故障排查求助

多线程生产者-消费者问题的错误排查与修复

你遇到的是典型的生产者-消费者(Writer-Reader)同步问题,代码在互斥锁使用、条件变量逻辑、缓冲区管理以及文件读取的线程安全上都踩了不少坑,才导致reader输出0、writer重复读取最后一个整数、线程无法正常终止这些问题。下面我一步步帮你拆解问题并给出修复方案:

一、核心错误点梳理

1. 同步原语未初始化+逻辑完全混乱

  • 你只声明了互斥锁pthread_mutex_t m;,但完全没调用pthread_mutex_init(&m, NULL);初始化,这会导致锁的行为完全不可控,线程随便乱闯临界区。
  • 用cnt变量控制线程等待的逻辑完全失效:writer里刚把cnt加1就立刻设为0,根本起不到“缓冲区满/空”的同步作用,多个writer线程能同时写入缓冲区,直接导致数据覆盖或重复读取。
  • reader的等待条件while (cnt == -1)完全不合理,初始cnt是0,这个条件永远不成立,reader不会等待直接输出缓冲区的初始值(全0)。

2. 文件读取的线程安全问题

多个writer线程直接共享fp调用fscanf,但fscanf不是线程安全函数——多个线程同时操作文件指针时,会出现数据竞争,这就是为什么两个writer都读到了最后一个整数的原因:文件指针的位置被多个线程同时修改,导致重复读取。

3. 缓冲区操作逻辑错误

  • 全局变量z用来记录写入位置,但没有和缓冲区满的判断绑定:你只在z==5时重置位置,但没判断缓冲区是否真的满了就继续写入,多个writer同时操作z还会导致位置混乱、数据覆盖。
  • reader直接循环输出buf的前5个元素,不管缓冲区里实际有没有有效数据,也没标记哪些数据已经被读取,自然会输出未写入的0。

4. 线程终止逻辑失效

单个writer读完文件后设置flag = -1,但另一个writer可能还在运行,而且flag没有被同步保护,reader无法准确感知所有writer是否完成,导致无限循环。

二、修复后的完整代码

下面是符合需求的修正代码,保留了2个writer、2个reader、缓冲区大小5的设定,保证所有数据被正确读写,线程能正常终止:

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

#define BUF_SIZE 5

// 同步原语:互斥锁+两个条件变量(分别对应缓冲区不满/不空)
pthread_mutex_t m;
pthread_cond_t cond_not_full;
pthread_cond_t cond_not_empty;

// 共享资源
FILE *fp;
int buf[BUF_SIZE];
int write_idx = 0;  // 缓冲区写入位置(循环队列)
int read_idx = 0;   // 缓冲区读取位置(循环队列)
int buf_count = 0;  // 缓冲区中有效数据的数量
int all_writer_done = 0;  // 标记所有writer是否完成

// 线程函数声明
void *reader(void *parm);
void *writer(void *parm);

// 线程安全的文件读取函数:加锁保护文件指针操作
int safe_fscanf(int *val) {
    pthread_mutex_lock(&m);
    int ret = fscanf(fp, "%d", val);
    pthread_mutex_unlock(&m);
    return ret;
}

int main(int argc, char *argv[]) {
    pthread_t rid[2];
    pthread_t wid[2];

    // 初始化同步原语
    pthread_mutex_init(&m, NULL);
    pthread_cond_init(&cond_not_full, NULL);
    pthread_cond_init(&cond_not_empty, NULL);

    // 打开文件
    if ((fp = fopen("/home/bhagi/2/shared_data.txt", "r")) == NULL) {
        fprintf(stderr, "Couldn't open the file\n");
        exit(EXIT_FAILURE);
    }

    // 创建writer线程
    for (int i = 0; i < 2; i++) {
        printf("Writer %d Starts \n", i + 1);
        pthread_create(&wid[i], NULL, writer, (void *)(long)i);
    }

    // 创建reader线程
    for (int i = 0; i < 2; i++) {
        printf("Reader %d Starts \n", i + 1);
        pthread_create(&rid[i], NULL, reader, (void *)(long)i);
    }

    // 等待所有writer完成
    for (int i = 0; i < 2; i++) {
        pthread_join(wid[i], NULL);
    }

    // 所有writer完成后,通知所有reader退出
    pthread_mutex_lock(&m);
    all_writer_done = 1;
    pthread_cond_broadcast(&cond_not_empty);
    pthread_mutex_unlock(&m);

    // 等待所有reader完成
    for (int i = 0; i < 2; i++) {
        pthread_join(rid[i], NULL);
    }

    // 清理资源
    fclose(fp);
    pthread_mutex_destroy(&m);
    pthread_cond_destroy(&cond_not_full);
    pthread_cond_destroy(&cond_not_empty);

    return 0;
}

void *reader(void *parm) {
    int id = (int)(long)parm;
    int read_count = 0;

    while (1) {
        pthread_mutex_lock(&m);

        // 等待:缓冲区空且writer未完成时阻塞
        while (buf_count == 0 && !all_writer_done) {
            printf("Reader %d waits for data\n", id + 1);
            pthread_cond_wait(&cond_not_empty, &m);
        }

        // 退出条件:缓冲区空且所有writer已完成
        if (buf_count == 0 && all_writer_done) {
            pthread_mutex_unlock(&m);
            break;
        }

        // 读取缓冲区数据(循环队列)
        int val = buf[read_idx];
        read_idx = (read_idx + 1) % BUF_SIZE;
        buf_count--;
        read_count++;

        printf("Reader %d read: %d (total read: %d)\n", id + 1, val, read_count);

        // 通知writer:缓冲区有空位了
        pthread_cond_signal(&cond_not_full);
        pthread_mutex_unlock(&m);

        sleep(1);  // 模拟处理时间,可按需删除
    }

    printf("Reader %d exits, total read: %d\n", id + 1, read_count);
    return NULL;
}

void *writer(void *parm) {
    int id = (int)(long)parm;
    int write_count = 0;
    int val;

    while (safe_fscanf(&val) == 1) {
        pthread_mutex_lock(&m);

        // 等待:缓冲区满时阻塞
        while (buf_count == BUF_SIZE) {
            printf("Writer %d waits for empty slot\n", id + 1);
            pthread_cond_wait(&cond_not_full, &m);
        }

        // 写入缓冲区(循环队列)
        buf[write_idx] = val;
        write_idx = (write_idx + 1) % BUF_SIZE;
        buf_count++;
        write_count++;

        printf("Writer %d wrote: %d (total written: %d)\n", id + 1, val, write_count);

        // 通知reader:缓冲区有数据了
        pthread_cond_signal(&cond_not_empty);
        pthread_mutex_unlock(&m);

        sleep(1);  // 模拟写入时间,可按需删除
    }

    printf("Writer %d exits, total written: %d\n", id + 1, write_count);
    return NULL;
}

三、关键修复说明

  1. 正确初始化同步原语:补全了互斥锁和条件变量的初始化,避免未定义行为。
  2. 线程安全的文件读取:用safe_fscanf加锁保护文件操作,防止多个writer同时修改文件指针导致数据竞争。
  3. 循环队列管理缓冲区:用write_idx、read_idx和buf_count实现循环队列,准确记录缓冲区的使用情况,避免数据覆盖或重复读取。
  4. 合理的同步逻辑:
    • Writer在缓冲区满时等待,写入后通知reader;
    • Reader在缓冲区空且writer未完成时等待,读取后通知writer;
    • 所有writer完成后,通过all_writer_done标记广播通知所有reader退出。
  5. 清晰的线程终止逻辑:确保reader能准确感知所有writer的状态,不会无限等待。

内容的提问来源于stack exchange,提问作者Doctor.008

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:57:12