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

基于信号量的C语言生产者-消费者同步问题优化咨询

Fixing Consumer Thread Synchronization in Producer-Consumer (No More sleep()!)

Great question! Using sleep() for thread synchronization is a classic hack—it might work on your machine right now, but it's unreliable. Timing can vary across systems, under load, or even between runs, leading to race conditions or hangs. Let's fix this properly and clean up the code along the way.

Root Cause of the Problem

The issue happens because when you create consumer threads, the main thread might not have finished setting up critical resources yet:

  • The semaphores for empty slots, filled slots, or mutex locks might not be fully initialized.
  • The shared buffer might not be allocated or ready for access.
  • Producers might not have started adding items to the buffer, leaving consumers stuck waiting (or worse, accessing uninitialized memory).

Instead of hoping sleep() gives the main thread enough time, we use a dedicated synchronization primitive to ensure consumers only start working after all setup is done.

Efficient Synchronization Solution: Initialization Semaphore

The simplest fix is to add an initialization semaphore that blocks consumers until the main thread has finished all setup steps. Here's how it works:

  1. Create a semaphore initialized to 0 (meaning all consumers will block on it initially).
  2. After the main thread finishes initializing buffers, semaphores, and starting producers, post to this semaphore once per consumer.
  3. Each consumer thread starts by waiting on this semaphore—they won't proceed until the main thread gives the green light.

Optimized Code Example

Here's a cleaned-up, fixed version of your producer-consumer implementation:

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

// Define constants for buffer size and thread counts
#define BUFFER_SIZE 5
#define NUM_PRODUCERS 2
#define NUM_CONSUMERS 2

// Struct to hold shared resources (avoids messy global variables!)
typedef struct {
    int buffer[BUFFER_SIZE];
    int in;          // Next position to add an item
    int out;         // Next position to remove an item
    sem_t empty;     // Tracks empty slots in the buffer
    sem_t full;      // Tracks filled slots in the buffer
    sem_t init_done; // Blocks consumers until setup is complete
    pthread_mutex_t mutex; // Protects buffer access from race conditions
    int stop;        // Flag to tell threads to exit gracefully
} SharedData;

// Producer thread function
void* producer(void* arg) {
    SharedData* data = (SharedData*)arg;
    int item;

    while (!data->stop) {
        // Simulate producing an item
        item = rand() % 100;
        printf("Producer %lu produced: %d\n", pthread_self(), item);

        // Wait for an empty slot in the buffer
        if (sem_wait(&data->empty) != 0) {
            perror("sem_wait(empty) failed");
            break;
        }

        // Lock the buffer before modifying it
        pthread_mutex_lock(&data->mutex);
        data->buffer[data->in] = item;
        data->in = (data->in + 1) % BUFFER_SIZE;
        pthread_mutex_unlock(&data->mutex);

        // Signal that a slot is now filled
        sem_post(&data->full);

        // Simulate variable production time
        sleep(rand() % 2);
    }

    printf("Producer %lu exiting\n", pthread_self());
    return NULL;
}

// Consumer thread function
void* consumer(void* arg) {
    SharedData* data = (SharedData*)arg;
    int item;

    // Wait until main thread finishes all setup steps
    if (sem_wait(&data->init_done) != 0) {
        perror("sem_wait(init_done) failed");
        return NULL;
    }

    while (!data->stop) {
        // Wait for a filled slot in the buffer
        if (sem_wait(&data->full) != 0) {
            perror("sem_wait(full) failed");
            break;
        }

        // Lock the buffer before accessing it
        pthread_mutex_lock(&data->mutex);
        item = data->buffer[data->out];
        data->out = (data->out + 1) % BUFFER_SIZE;
        pthread_mutex_unlock(&data->mutex);

        // Signal that a slot is now empty
        sem_post(&data->empty);

        // Simulate consuming the item
        printf("Consumer %lu consumed: %d\n", pthread_self(), item);

        // Simulate variable consumption time
        sleep(rand() % 2);
    }

    printf("Consumer %lu exiting\n", pthread_self());
    return NULL;
}

int main() {
    SharedData data;
    pthread_t producers[NUM_PRODUCERS];
    pthread_t consumers[NUM_CONSUMERS];
    int i;

    // Initialize shared state
    data.in = 0;
    data.out = 0;
    data.stop = 0;

    // Initialize semaphores
    if (sem_init(&data.empty, 0, BUFFER_SIZE) != 0 ||
        sem_init(&data.full, 0, 0) != 0 ||
        sem_init(&data.init_done, 0, 0) != 0) {
        perror("sem_init failed");
        exit(EXIT_FAILURE);
    }

    // Initialize mutex lock
    if (pthread_mutex_init(&data.mutex, NULL) != 0) {
        perror("pthread_mutex_init failed");
        exit(EXIT_FAILURE);
    }

    // Create producer threads
    for (i = 0; i < NUM_PRODUCERS; i++) {
        if (pthread_create(&producers[i], NULL, producer, &data) != 0) {
            perror("pthread_create(producer) failed");
            data.stop = 1;
            exit(EXIT_FAILURE);
        }
    }

    // Create consumer threads
    for (i = 0; i < NUM_CONSUMERS; i++) {
        if (pthread_create(&consumers[i], NULL, consumer, &data) != 0) {
            perror("pthread_create(consumer) failed");
            data.stop = 1;
            exit(EXIT_FAILURE);
        }
    }

    // All setup complete! Let consumers start working
    for (i = 0; i < NUM_CONSUMERS; i++) {
        sem_post(&data.init_done);
    }

    // Let the system run for a demo period
    sleep(10);

    // Signal all threads to exit gracefully
    data.stop = 1;

    // Wake up any waiting producers/consumers so they can exit
    for (i = 0; i < NUM_PRODUCERS; i++) {
        sem_post(&data.empty);
    }
    for (i = 0; i < NUM_CONSUMERS; i++) {
        sem_post(&data.full);
    }

    // Wait for all threads to finish
    for (i = 0; i < NUM_PRODUCERS; i++) {
        pthread_join(producers[i], NULL);
    }
    for (i = 0; i < NUM_CONSUMERS; i++) {
        pthread_join(consumers[i], NULL);
    }

    // Clean up resources to avoid leaks
    sem_destroy(&data.empty);
    sem_destroy(&data.full);
    sem_destroy(&data.init_done);
    pthread_mutex_destroy(&data.mutex);

    printf("All threads exited. Cleanup complete.\n");
    return 0;
}

Key Optimizations & Fixes

Let's break down the improvements:

  • No global variables: All shared state is wrapped in a SharedData struct, making the code modular and eliminating unintended side effects from global scope.
  • Reliable initialization sync: The init_done semaphore ensures consumers only start after all setup is complete—no more guessing with sleep().
  • Error checking: Every system call (semaphore operations, thread creation, mutex init) has error handling, making debugging far easier.
  • Graceful thread exit: The stop flag lets threads exit cleanly instead of being force-killed. We also wake up waiting threads with semaphore posts so they don't hang indefinitely.
  • Proper resource cleanup: All semaphores and mutexes are destroyed after use, preventing resource leaks.

Alternative: Using Condition Variables

If you prefer mutexes and condition variables over semaphores for initialization, you could add a pthread_cond_t init_cond and boolean is_initialized flag to the SharedData struct. The main thread would set is_initialized = true and signal the condition variable, while consumers wait on it until the flag is set. The semaphore approach is simpler for this specific use case, though.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 08:29:16