基于pthread条件变量的两类多生产者多消费者同步问题求解
Hey there! Let's tackle your two pthread condition variable questions step by step—this is classic synchronization territory, so I'll break it down with clear logic and working code examples to make it stick.
1. 如何使用pthread条件变量解决多生产者多消费者的同步问题?
First, let's recap the core tools you need here:
- A mutex (
pthread_mutex_t) to protect access to the shared queue (since multiple threads will be reading/writing to it). - Two condition variables (
pthread_cond_t): one to signal when the queue isn't empty (for consumers waiting to take data) and another to signal when the queue isn't full (for producers waiting to add data). - A bounded or unbounded queue (we'll use a bounded one here because it's more realistic for resource limits).
Key Rules to Follow:
- Always lock the mutex before checking queue state or modifying the queue.
- Use a
whileloop (notif) to check queue conditions—this handles spurious wakeups (when a condition variable wakes up without the actual condition being met). - Signal/broadcast the condition variable after modifying the queue to wake up waiting threads.
Working Code Example:
#include <pthread.h> #include <stdio.h> #include <stdlib.h> #include <unistd.h> #define QUEUE_SIZE 10 #define NUM_PRODUCERS 2 #define NUM_CONSUMERS 2 typedef struct { int data[QUEUE_SIZE]; int front; int rear; int count; pthread_mutex_t mutex; pthread_cond_t not_full; pthread_cond_t not_empty; } BoundedQueue; void init_queue(BoundedQueue *q) { q->front = 0; q->rear = 0; q->count = 0; pthread_mutex_init(&q->mutex, NULL); pthread_cond_init(&q->not_full, NULL); pthread_cond_init(&q->not_empty, NULL); } void enqueue(BoundedQueue *q, int item) { pthread_mutex_lock(&q->mutex); // Wait until queue has space while (q->count == QUEUE_SIZE) { printf("Queue full, producer waiting...\n"); pthread_cond_wait(&q->not_full, &q->mutex); } q->data[q->rear] = item; q->rear = (q->rear + 1) % QUEUE_SIZE; q->count++; printf("Produced %d, queue count: %d\n", item, q->count); // Wake up a waiting consumer pthread_cond_signal(&q->not_empty); pthread_mutex_unlock(&q->mutex); } int dequeue(BoundedQueue *q) { pthread_mutex_lock(&q->mutex); // Wait until queue has items while (q->count == 0) { printf("Queue empty, consumer waiting...\n"); pthread_cond_wait(&q->not_empty, &q->mutex); } int item = q->data[q->front]; q->front = (q->front + 1) % QUEUE_SIZE; q->count--; printf("Consumed %d, queue count: %d\n", item, q->count); // Wake up a waiting producer pthread_cond_signal(&q->not_full); pthread_mutex_unlock(&q->mutex); return item; } void *producer(void *arg) { BoundedQueue *q = (BoundedQueue *)arg; for (int i = 0; i < 5; i++) { enqueue(q, i + 100); sleep(rand() % 2); // Simulate work time } return NULL; } void *consumer(void *arg) { BoundedQueue *q = (BoundedQueue *)arg; for (int i = 0; i < 5; i++) { int item = dequeue(q); sleep(rand() % 3); // Simulate work time } return NULL; } int main() { BoundedQueue q; init_queue(&q); pthread_t prod_threads[NUM_PRODUCERS]; pthread_t cons_threads[NUM_CONSUMERS]; for (int i = 0; i < NUM_PRODUCERS; i++) { pthread_create(&prod_threads[i], NULL, producer, &q); } for (int i = 0; i < NUM_CONSUMERS; i++) { pthread_create(&cons_threads[i], NULL, consumer, &q); } for (int i = 0; i < NUM_PRODUCERS; i++) { pthread_join(prod_threads[i], NULL); } for (int i = 0; i < NUM_CONSUMERS; i++) { pthread_join(cons_threads[i], NULL); } pthread_mutex_destroy(&q->mutex); pthread_cond_destroy(&q->not_full); pthread_cond_destroy(&q->not_empty); return 0; }
Quick Explanation:
- Producers add items to the queue; if it's full, they wait on
not_fulluntil a consumer frees up space. - Consumers take items from the queue; if it's empty, they wait on
not_emptyuntil a producer adds data. - The
whileloop ensures that even if a thread is woken up spuriously, it rechecks the queue state before proceeding.
2. 现有三个生产者,每个拥有一个专属队列;两个消费者需同时从指定专属队列与公共队列消费消息,如何通过pthread_cond实现同步控制?
This is a more nuanced scenario—let's clarify the setup first:
- 3 producers (P1, P2, P3) each write only to their own exclusive queue (Q1, Q2, Q3).
- 2 consumers (C1, C2): C1 reads from Q1 + public queue Q4; C2 reads from Q2 + public queue Q4 (adjust if your "specified exclusive queue" means something else, but this is a common use case).
Core Approach:
- Each queue (exclusive + public) gets its own mutex to protect queue operations.
- A global condition variable to wake up consumers whenever any of their target queues gets new data. This avoids having consumers poll the queues repeatedly.
- Consumers first check their exclusive queue, then the public queue—if both are empty, they wait on the global condition variable.
Working Code Example:
#include <pthread.h> #include <stdio.h> #include <stdlib.h> #include <unistd.h> #define QUEUE_CAPACITY 5 // Generic queue structure typedef struct { int items[QUEUE_CAPACITY]; int front; int rear; int count; pthread_mutex_t mutex; } Queue; // Global sync structure to wake consumers typedef struct { pthread_cond_t consumer_cond; pthread_mutex_t global_mutex; } GlobalSync; Queue q1, q2, q3, q_public; GlobalSync sync; void init_queue(Queue *q) { q->front = q->rear = q->count = 0; pthread_mutex_init(&q->mutex, NULL); } void init_global_sync(GlobalSync *s) { pthread_cond_init(&s->consumer_cond, NULL); pthread_mutex_init(&s->global_mutex, NULL); } // Enqueue item and wake consumers void enqueue(Queue *q, int item, const char *queue_name) { pthread_mutex_lock(&q->mutex); if (q->count == QUEUE_CAPACITY) { printf("Producer: %s is full, dropping item %d\n", queue_name, item); pthread_mutex_unlock(&q->mutex); return; } q->items[q->rear] = item; q->rear = (q->rear + 1) % QUEUE_CAPACITY; q->count++; printf("Produced to %s: %d, count: %d\n", queue_name, item, q->count); pthread_mutex_unlock(&q->mutex); // Wake all waiting consumers pthread_mutex_lock(&sync.global_mutex); pthread_cond_broadcast(&sync.consumer_cond); pthread_mutex_unlock(&sync.global_mutex); } // Dequeue item (returns 1 if successful, 0 if empty) int dequeue(Queue *q, int *item, const char *queue_name) { pthread_mutex_lock(&q->mutex); if (q->count == 0) { pthread_mutex_unlock(&q->mutex); return 0; } *item = q->items[q->front]; q->front = (q->front + 1) % QUEUE_CAPACITY; q->count--; printf("Consumed from %s: %d, count: %d\n", queue_name, *item, q->count); pthread_mutex_unlock(&q->mutex); return 1; } // Producer 1: writes only to Q1 void *producer1(void *arg) { for (int i = 0; i < 4; i++) { enqueue(&q1, i + 10, "Q1"); sleep(rand() % 2); } return NULL; } // Producer 2: writes only to Q2 void *producer2(void *arg) { for (int i = 0; i < 4; i++) { enqueue(&q2, i + 20, "Q2"); sleep(rand() % 2); } return NULL; } // Producer 3: writes only to Q3 void *producer3(void *arg) { for (int i = 0; i < 4; i++) { enqueue(&q3, i + 30, "Q3"); sleep(rand() % 2); } return NULL; } // Optional: Producer for public queue void *public_producer(void *arg) { for (int i = 0; i < 5; i++) { enqueue(&q_public, i + 100, "Q_PUBLIC"); sleep(rand() % 3); } return NULL; } // Consumer 1: reads Q1 + public queue void *consumer1(void *arg) { int item; while (1) { // Try exclusive queue first if (dequeue(&q1, &item, "Q1")) { sleep(rand() % 2); // Simulate processing continue; } // Try public queue next if (dequeue(&q_public, &item, "Q_PUBLIC")) { sleep(rand() % 2); continue; } // Wait for new data pthread_mutex_lock(&sync.global_mutex); printf("Consumer1: No items available, waiting...\n"); pthread_cond_wait(&sync.consumer_cond, &sync.global_mutex); pthread_mutex_unlock(&sync.global_mutex); } return NULL; } // Consumer 2: reads Q2 + public queue void *consumer2(void *arg) { int item; while (1) { if (dequeue(&q2, &item, "Q2")) { sleep(rand() % 2); continue; } if (dequeue(&q_public, &item, "Q_PUBLIC")) { sleep(rand() % 2); continue; } pthread_mutex_lock(&sync.global_mutex); printf("Consumer2: No items available, waiting...\n"); pthread_cond_wait(&sync.consumer_cond, &sync.global_mutex); pthread_mutex_unlock(&sync.global_mutex); } return NULL; } int main() { // Initialize all queues and global sync init_queue(&q1); init_queue(&q2); init_queue(&q3); init_queue(&q_public); init_global_sync(&sync); pthread_t p1, p2, p3, p_public, c1, c2; pthread_create(&p1, NULL, producer1, NULL); pthread_create(&p2, NULL, producer2, NULL); pthread_create(&p3, NULL, producer3, NULL); pthread_create(&p_public, NULL, public_producer, NULL); pthread_create(&c1, NULL, consumer1, NULL); pthread_create(&c2, NULL, consumer2, NULL); // Wait for producers to finish pthread_join(p1, NULL); pthread_join(p2, NULL); pthread_join(p3, NULL); pthread_join(p_public, NULL); // Give consumers time to process remaining items sleep(5); printf("Main thread exiting...\n"); return 0; }
Key Notes:
- We use
pthread_cond_broadcastinstead ofsignalbecause multiple consumers might be waiting for different queues—broadcast ensures all waiting consumers wake up and check their target queues. - Each queue's mutex keeps operations atomic, so we don't have race conditions when adding/removing items.
- If you need producers to wait for their exclusive queues to have space, you can add a
not_fullcondition variable per queue, similar to the first example.
内容的提问来源于stack exchange,提问作者wangsquirrel
相关产品推荐
相关产品推荐

