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

Java实现K个消费者线程被动等待的并行任务执行问询

Java K个消费者线程被动等待执行任务的实现方案

核心思路

用阻塞队列(如LinkedBlockingQueue)实现被动等待:队列的take()方法会在队列空时自动阻塞线程,不会像轮询那样占用CPU资源;任务执行完成后自动取下一个,无需主动轮询。针对不同场景,分两种实现方式:


场景1:一次性批处理任务(执行完所有任务后线程退出)

如果是一次性提交所有任务,执行完成后让消费者线程自动退出,需要给每个消费者发送终止标记(比如null),避免线程一直阻塞在take()。

消费者线程代码

class Consumer implements Runnable {
    private final BlockingQueue<Runnable> taskQueue;

    public Consumer(BlockingQueue<Runnable> taskQueue) {
        this.taskQueue = taskQueue;
    }

    @Override
    public void run() {
        try {
            while (true) {
                // 被动等待:队列空时自动阻塞,直到有任务或中断
                Runnable task = taskQueue.take();
                
                // 拿到终止标记,退出循环
                if (task == null) {
                    break;
                }
                
                // 执行任务,完成后自动回到take()取下一个
                task.run();
            }
        } catch (InterruptedException e) {
            // 线程被中断时,恢复中断状态并退出
            Thread.currentThread().interrupt();
        }
    }
}

任务提交与执行代码

public class BatchTaskExecutor {
    public static void execute(int threadCount, List<Runnable> tasks) {
        BlockingQueue<Runnable> taskQueue = new LinkedBlockingQueue<>();

        // 启动指定数量的消费者线程
        for (int i = 0; i < threadCount; i++) {
            new Thread(new Consumer(taskQueue)).start();
        }

        // 提交所有任务到队列
        try {
            for (Runnable task : tasks) {
                taskQueue.put(task);
            }

            // 提交与线程数相同的终止标记,确保每个消费者都能收到并退出
            for (int i = 0; i < threadCount; i++) {
                taskQueue.put(null);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

场景2:持续接收任务(线程长期运行)

如果需要线程一直运行,随时接收新提交的任务,无需终止标记,线程会一直阻塞在take(),直到有新任务进来立刻执行,或线程被中断。

持久化消费者线程代码

class PersistentConsumer implements Runnable {
    private final BlockingQueue<Runnable> taskQueue;

    public PersistentConsumer(BlockingQueue<Runnable> taskQueue) {
        this.taskQueue = taskQueue;
    }

    @Override
    public void run() {
        try {
            // 线程未被中断时,持续等待并执行任务
            while (!Thread.currentThread().isInterrupted()) {
                Runnable task = taskQueue.take();
                task.run();
            }
        } catch (InterruptedException e) {
            // 恢复中断状态,上层可处理
            Thread.currentThread().interrupt();
        }
    }
}

使用示例

public class PersistentTaskExecutor {
    public static void main(String[] args) {
        BlockingQueue<Runnable> taskQueue = new LinkedBlockingQueue<>();
        int threadCount = 3;

        // 启动持久化消费者线程
        for (int i = 0; i < threadCount; i++) {
            new Thread(new PersistentConsumer(taskQueue)).start();
        }

        // 随时提交新任务,消费者会立即执行
        taskQueue.put(() -> System.out.println("任务1执行中"));
        taskQueue.put(() -> System.out.println("任务2执行中"));
        
        // 若需停止线程,调用interrupt()即可
        // consumerThread.interrupt();
    }
}

解决你遇到的问题

  1. 主动等待改被动等待:用taskQueue.take()替代轮询逻辑,队列空时线程自动阻塞,完全避免主动轮询的CPU浪费。
  2. 有效任务数少于K时的循环问题:take()会阻塞线程,不会无意义循环,直到有任务或终止信号。
  3. 调用take()仅执行K个任务就停止:如果是一次性任务场景,是因为任务执行完后线程阻塞在take(),但主线程结束后消费者线程仍在后台运行;添加终止标记后,线程会执行完最后一个任务后主动退出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 07:45:01