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

C语言中如何终止阻塞在空/满数据结构上的生产者/消费者线程?

问题描述

假设有若干消费者线程从阻塞队列(BlockingQueue)中无限循环读取元素,生产者线程向该队列生产一定数量元素后退出。消费者不知道元素总数,因此持续循环,在所有元素生产并消费完成后,消费者线程会因等待更多元素而阻塞。目标是编写规范代码实现消费者线程的终止。

在Java中有一种抽象性良好的实现方案:通过Thread.interrupt()和异常机制实现,无需访问数据结构内部实现,仅需调用统一的Thread.interrupt()即可终止线程。以下为1个消费者和1个生产者的示例代码(久未使用Java,若泛型等存在错误敬请谅解):

public class Consume implements Runnable {

    private BlockingQueue<T> queue;

    public Consume(BlockingQueue<T> queue){
        this.queue = queue;
    }

    public void run() {
        try {
            while(true) {
                T elem;
                elem = queue.take(); //Blocking
                //Do something with the element
            }
        } catch(InterruptedException e) {
            //Do nothing and quit 
        }
    }
}

public class Produce implements Runnable {

    private BlockingQueue<T> queue;

    public Produce(BlockingQueue<T> queue) {
        this.queue = queue;
    }

    public void run() {
        //Insert a random number of element in the queue,  and then quit
    }
}

public static void main(String[] args) {
    BlockingQueue<Integer> queue = new ArrayBlockingQueue<>(capacity);
    Thread consumer = new Thread(new Consume(queue));
    Thread producer = new Thread(new Produce(queue));

    consumer.start();
    producer.start();
    producer.join();
    consumer.interrupt();

    return 0;
}

该方案的优势在于无需访问数据结构内部实现,且对各类数据结构均可使用同一函数终止线程。

在C语言中使用pthreads时,可通过pthread_cond_signal/pthread_cond_broadcast结合标记位实现,但需要访问数据结构的内部实现。

请问:

  • 是否存在更通用的解决方案?
  • 若没有,能否在C语言中实现类似Java的方案?

我发现macOS中,等待pthread_cond_wait的线程收到系统信号会返回,基于此编写了类似Java的方案,但Linux中pthread_cond_wait不会因信号返回,而我的代码需要兼容两种系统。


解决方案

一、通用跨平台终止方案(无侵入数据结构)

C语言没有Java那种线程级统一中断机制,但可以通过自定义线程终止标记+安全条件变量等待逻辑实现接近的抽象效果,且无需修改阻塞队列内部实现。核心是给每个消费者线程绑定终止标记,将标记检查与条件等待结合。

实现步骤

  1. 定义线程上下文结构体,封装队列相关资源与终止标记:
#include <pthread.h>
#include <stdbool.h>

// 消费者线程上下文
typedef struct {
    void* queue;                  // 阻塞队列抽象指针(无需知晓内部结构)
    pthread_mutex_t* queue_mutex; // 队列的互斥锁(由队列提供)
    pthread_cond_t* queue_cond;   // 队列的条件变量(由队列提供)
    bool should_terminate;        // 终止标记
    pthread_mutex_t terminate_mutex; // 保护终止标记的互斥锁
} ConsumerCtx;
  1. 消费者线程逻辑:每次等待元素前先检查终止标记,条件变量等待时处理超时或中断,确保能响应终止信号:
// 假设队列提供的辅助函数:检查队列是否为空、取出元素
bool queue_is_empty(void* queue);
void* queue_take(void* queue);
// 元素处理函数
void process_element(void* elem);

void* consumer_thread(void* arg) {
    ConsumerCtx* ctx = (ConsumerCtx*)arg;
    void* elem;
    bool terminate = false;

    while (true) {
        // 检查终止标记
        pthread_mutex_lock(&ctx->terminate_mutex);
        terminate = ctx->should_terminate;
        pthread_mutex_unlock(&ctx->terminate_mutex);
        if (terminate) break;

        pthread_mutex_lock(ctx->queue_mutex);
        // 等待队列非空或终止标记触发
        while (queue_is_empty(ctx->queue) && !terminate) {
            // 使用超时等待避免永久阻塞,同时兼容信号中断
            struct timespec ts;
            clock_gettime(CLOCK_REALTIME, &ts);
            ts.tv_sec += 1; // 1秒超时
            int ret = pthread_cond_timedwait(ctx->queue_cond, ctx->queue_mutex, &ts);
            
            if (ret == ETIMEDOUT || ret == EINTR) {
                // 超时或被信号中断后,重新检查终止标记
                pthread_mutex_lock(&ctx->terminate_mutex);
                terminate = ctx->should_terminate;
                pthread_mutex_unlock(&ctx->terminate_mutex);
                continue;
            }
        }

        if (terminate) {
            pthread_mutex_unlock(ctx->queue_mutex);
            break;
        }

        // 取出并处理元素
        elem = queue_take(ctx->queue);
        pthread_mutex_unlock(ctx->queue_mutex);
        process_element(elem);
    }

    return NULL;
}
  1. 统一终止接口:设置标记并唤醒等待线程,确保消费者能退出阻塞:
void terminate_consumer(ConsumerCtx* ctx) {
    pthread_mutex_lock(&ctx->terminate_mutex);
    ctx->should_terminate = true;
    pthread_mutex_unlock(&ctx->terminate_mutex);
    // 广播唤醒所有等待该条件变量的线程
    pthread_cond_broadcast(ctx->queue_cond);
}

该方案优势:

  • 无需修改阻塞队列内部,仅需队列暴露互斥锁、条件变量和空状态检查接口;
  • 完全跨平台,不受macOS/Linux信号处理差异影响;
  • 抽象程度接近Java的Thread.interrupt(),通过统一函数终止线程。

二、模拟Java风格的信号中断方案(兼容跨平台)

若要模拟Java的线程中断行为,需处理两大系统的信号差异:

  1. 给每个消费者绑定专属信号(如SIGUSR1),用线程局部存储标记中断状态;
  2. 处理pthread_cond_wait的返回值:Linux下返回EINTR,macOS下直接唤醒,都需检查中断标记。

示例代码

#include <pthread.h>
#include <signal.h>
#include <stdbool.h>

// 线程局部中断标记
__thread bool thread_interrupted = false;

void sigusr1_handler(int sig) {
    thread_interrupted = true;
}

// 假设阻塞队列的取元素函数内部调用pthread_cond_wait
int queue_take_blocking(void* queue, void** elem);
void process_element(void* elem);

void* consumer_thread(void* arg) {
    void* queue = arg;
    void* elem;

    // 为当前线程注册SIGUSR1信号处理
    struct sigaction sa;
    sa.sa_handler = sigusr1_handler;
    sigemptyset(&sa.sa_mask);
    sa.sa_flags = 0;
    sigaction(SIGUSR1, &sa, NULL);

    while (true) {
        if (thread_interrupted) break;

        // 尝试取元素,处理中断
        int ret = queue_take_blocking(queue, &elem);
        if (ret == EINTR || thread_interrupted) break;

        process_element(elem);
    }

    return NULL;
}

// 发送信号终止线程
void interrupt_thread(pthread_t tid) {
    pthread_kill(tid, SIGUSR1);
}

该方案局限性:

  • 要求队列的阻塞取函数不透明重试EINTR,否则无法响应中断;
  • 信号处理可能引发竞态,需谨慎维护线程局部标记;
  • 需额外测试跨平台信号行为差异。

三、推荐方案

优先选择自定义终止标记+条件变量超时等待的方案,原因:

  • 无信号处理的竞态风险;
  • 完全跨平台,不受系统条件变量行为差异影响;
  • 逻辑清晰,易于维护和调试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 23:33:08