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

多线程生产者-消费者问题:超范围生产及线程退出异常求助

Pthread Producer-Consumer: Producers Generate Items Beyond Limit and Threads Fail to Exit

I've got a pthread-based producer-consumer implementation that's supposed to produce and consume items within the range specified by the total_items parameter (e.g., 0-39 when total_items=40). However, two critical issues are occurring:

  1. Producers are generating items beyond the total_items limit (like 40 and higher in the example).
  2. Both producer and consumer threads fail to exit properly when their exit conditions are met.

Here's the problematic code:

#include <stdio.h>
#include <stdlib.h>
#include <sys/types.h>
#include <string.h>
#include <errno.h>
#include <signal.h>
#include <wait.h>
#include <pthread.h>

int item_to_produce, curr_buf_size;
int total_items, max_buf_size, num_workers, num_masters;
int consumed_items;
int *buffer;
pthread_mutex_t mutex;
pthread_cond_t has_data;
pthread_cond_t has_space;

void print_produced(int num, int master) {
    printf("Produced %d by master %d\n", num, master);
}

void print_consumed(int num, int worker) {
    printf("Consumed %d by worker %d\n", num, worker);
}

//consume items in buffer
void *consume_requests_loop(void *data) {
    int thread_id = *((int *)data);
    while(1) {
        pthread_mutex_lock(&mutex);
        if(consumed_items == total_items) {
            pthread_mutex_unlock(&mutex);
            break;
        }
        while(curr_buf_size == 0) {
            pthread_cond_wait(&has_data, &mutex);
        }
        print_consumed(buffer[(curr_buf_size--)-1], thread_id);
        consumed_items++;
        pthread_cond_signal(&has_space);
        pthread_mutex_unlock(&mutex);
    }
    return 0;
}

//produce items and place in buffer
void *generate_requests_loop(void *data) {
    int thread_id = *((int *)data);
    while(1) {
        pthread_mutex_lock(&mutex);
        if(item_to_produce == total_items) {
            pthread_mutex_unlock(&mutex);
            break;
        }
        while (curr_buf_size == max_buf_size) {
            pthread_cond_wait(&has_space, &mutex);
        }
        buffer[curr_buf_size++] = item_to_produce;
        print_produced(item_to_produce, thread_id);
        item_to_produce++;
        pthread_cond_signal(&has_data);
        pthread_mutex_unlock(&mutex);
    }
    return 0;
}

int main(int argc, char *argv[]) {
    int *master_thread_id;
    int *worker_thread_id;
    pthread_t *master_thread;
    pthread_t *worker_thread;

    item_to_produce = 0;
    curr_buf_size = 0;
    consumed_items = 0;
    int i;

    if (argc < 5) {
        printf("./master-worker #total_items #max_buf_size #num_workers #masters e.g. ./exe 10000 1000 4 3\n");
        exit(1);
    } else {
        num_masters = atoi(argv[4]);
        num_workers = atoi(argv[3]);
        total_items = atoi(argv[1]);
        max_buf_size = atoi(argv[2]);
    }

    buffer = (int *)malloc (sizeof(int) * max_buf_size);
    pthread_mutex_init(&mutex, NULL);
    pthread_cond_init(&has_space, NULL);
    pthread_cond_init(&has_data, NULL);

    //create master producer threads
    master_thread_id = (int *)malloc(sizeof(int) * num_masters);
    master_thread = (pthread_t *)malloc(sizeof(pthread_t) * num_masters);
    for (i = 0; i < num_masters; i++)
        master_thread_id[i] = i;
    for (i = 0; i < num_masters; i++)
        pthread_create(&master_thread[i], NULL, generate_requests_loop, (void *)&master_thread_id[i]);

    //create worker consumer threads
    worker_thread_id = (int *)malloc(sizeof(int) * num_workers);
    worker_thread = (pthread_t *)malloc(sizeof(pthread_t) * num_workers);
    for (i = 0; i < num_workers; i++)
        worker_thread_id[i] = i;
    for (i = 0 ; i < num_workers; i++)
        pthread_create(&worker_thread[i], NULL, consume_requests_loop, (void *)&worker_thread_id[i]);

    //wait for all threads to complete
    for (i = 0; i < num_masters; i++) {
        pthread_join(master_thread[i], NULL);
        printf("master %d joined\n", i);
    }

    for (i = 0; i < num_workers; i++) {
        pthread_join(worker_thread[i], NULL);
        printf("worker %d joined\n", i);
    }

    free(buffer);
    free(master_thread_id);
    free(master_thread);
    free(worker_thread_id);
    free(worker_thread);

    pthread_mutex_destroy(&mutex);
    pthread_cond_destroy(&has_data);
    pthread_cond_destroy(&has_space);

    return 0;
}

Root Causes of the Issues

1. Out-of-Bounds Production

The producer threads check the exit condition (item_to_produce == total_items) before entering the condition wait for buffer space. If a producer is blocked waiting for space, other producers may finish producing the last item, pushing item_to_produce to total_items. When the blocked producer is awakened, it skips the exit check and proceeds to produce an item using the now-invalid item_to_produce value.

2. Threads Stuck in Condition Waits

When all producers exit, some consumer threads may still be blocked waiting for has_data signals. Since no producers are left to send these signals, the consumers hang indefinitely. Similarly, producers could get stuck waiting for has_space if all consumers exit early (though less likely here).


Fixed Code

#include <stdio.h>
#include <stdlib.h>
#include <sys/types.h>
#include <string.h>
#include <errno.h>
#include <signal.h>
#include <wait.h>
#include <pthread.h>

int item_to_produce, curr_buf_size;
int total_items, max_buf_size, num_workers, num_masters;
int consumed_items;
int *buffer;
pthread_mutex_t mutex;
pthread_cond_t has_data;
pthread_cond_t has_space;

void print_produced(int num, int master) {
    printf("Produced %d by master %d\n", num, master);
}

void print_consumed(int num, int worker) {
    printf("Consumed %d by worker %d\n", num, worker);
}

//consume items in buffer
void *consume_requests_loop(void *data) {
    int thread_id = *((int *)data);
    while(1) {
        pthread_mutex_lock(&mutex);
        
        // Wait until buffer has data OR we've consumed all items
        while (curr_buf_size == 0 && consumed_items < total_items) {
            pthread_cond_wait(&has_data, &mutex);
        }
        
        // Check if we need to exit
        if(consumed_items == total_items) {
            pthread_mutex_unlock(&mutex);
            // Wake other consumers so they can check exit condition
            pthread_cond_broadcast(&has_data);
            break;
        }
        
        print_consumed(buffer[(curr_buf_size--)-1], thread_id);
        consumed_items++;
        pthread_cond_signal(&has_space);
        pthread_mutex_unlock(&mutex);
    }
    return 0;
}

//produce items and place in buffer
void *generate_requests_loop(void *data) {
    int thread_id = *((int *)data);
    while(1) {
        pthread_mutex_lock(&mutex);
        
        // Wait until buffer has space OR we've produced all items
        while (curr_buf_size == max_buf_size && item_to_produce < total_items) {
            pthread_cond_wait(&has_space, &mutex);
        }
        
        // Check if we need to exit
        if(item_to_produce == total_items) {
            pthread_mutex_unlock(&mutex);
            // Wake other producers so they can check exit condition
            pthread_cond_broadcast(&has_space);
            // Wake consumers in case they're waiting for remaining items
            pthread_cond_broadcast(&has_data);
            break;
        }
        
        buffer[curr_buf_size++] = item_to_produce;
        print_produced(item_to_produce, thread_id);
        item_to_produce++;
        pthread_cond_signal(&has_data);
        pthread_mutex_unlock(&mutex);
    }
    return 0;
}

int main(int argc, char *argv[]) {
    int *master_thread_id;
    int *worker_thread_id;
    pthread_t *master_thread;
    pthread_t *worker_thread;

    item_to_produce = 0;
    curr_buf_size = 0;
    consumed_items = 0;
    int i;

    if (argc < 5) {
        printf("./master-worker #total_items #max_buf_size #num_workers #masters e.g. ./exe 10000 1000 4 3\n");
        exit(1);
    } else {
        num_masters = atoi(argv[4]);
        num_workers = atoi(argv[3]);
        total_items = atoi(argv[1]);
        max_buf_size = atoi(argv[2]);
    }

    buffer = (int *)malloc (sizeof(int) * max_buf_size);
    pthread_mutex_init(&mutex, NULL);
    pthread_cond_init(&has_space, NULL);
    pthread_cond_init(&has_data, NULL);

    //create master producer threads
    master_thread_id = (int *)malloc(sizeof(int) * num_masters);
    master_thread = (pthread_t *)malloc(sizeof(pthread_t) * num_masters);
    for (i = 0; i < num_masters; i++)
        master_thread_id[i] = i;
    for (i = 0; i < num_masters; i++)
        pthread_create(&master_thread[i], NULL, generate_requests_loop, (void *)&master_thread_id[i]);

    //create worker consumer threads
    worker_thread_id = (int *)malloc(sizeof(int) * num_workers);
    worker_thread = (pthread_t *)malloc(sizeof(pthread_t) * num_workers);
    for (i = 0; i < num_workers; i++)
        worker_thread_id[i] = i;
    for (i = 0 ; i < num_workers; i++)
        pthread_create(&worker_thread[i], NULL, consume_requests_loop, (void *)&worker_thread_id[i]);

    //wait for all threads to complete
    for (i = 0; i < num_masters; i++) {
        pthread_join(master_thread[i], NULL);
        printf("master %d joined\n", i);
    }

    for (i = 0; i < num_workers; i++) {
        pthread_join(worker_thread[i], NULL);
        printf("worker %d joined\n", i);
    }

    free(buffer);
    free(master_thread_id);
    free(master_thread);
    free(worker_thread_id);
    free(worker_thread);

    pthread_mutex_destroy(&mutex);
    pthread_cond_destroy(&has_data);
    pthread_cond_destroy(&has_space);

    return 0;
}

Key Fixes Explained

For Producers:

  1. Moved Exit Check to After Condition Wait: The exit condition (item_to_produce == total_items) is checked after waking up from the has_space wait. This ensures we don't produce items once we've hit the total limit, even if we were blocked waiting for buffer space.
  2. Broadcast Signals on Exit: When a producer exits, it broadcasts has_space to wake other blocked producers, letting them check their own exit conditions. It also broadcasts has_data to ensure any waiting consumers wake up to process remaining items.
  3. Updated Wait Condition: The while loop for waiting now includes a check for whether we still need to produce items (item_to_produce < total_items), so we don't wait unnecessarily once production is complete.

For Consumers:

  1. Moved Exit Check to After Condition Wait: The exit condition (consumed_items == total_items) is checked after waking up from the has_data wait, preventing attempts to consume from an empty buffer once all items are processed.
  2. Broadcast Signals on Exit: When a consumer exits, it broadcasts has_data to wake other blocked consumers, allowing them to check their exit conditions and exit gracefully.
  3. Updated Wait Condition: The while loop for waiting includes a check for whether we still need to consume items (consumed_items < total_items), so we don't wait once all items are consumed.

These changes ensure no out-of-bounds production and all threads exit properly when their work is done.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:03:00