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

Java并行执行方法循环无法退出问题求助

并行任务执行后消费者线程无法退出的解决方案

问题分析

当前代码中,生产者线程将所有任务放入ArrayBlockingQueue后就结束,但消费者线程的while((r= taskQueue.take())!=null)循环会在队列空了之后一直阻塞在take()方法上——因为没有任何结束信号通知消费者任务已全部完成,导致主线程调用consumer.join()时永远等待。

解决思路

在生产者完成所有任务的提交后,向队列中放入与消费者线程数量相等的null值作为结束标记。每个消费者线程在拿到null时,就知道没有更多任务需要处理,从而退出循环。这种方式不需要主动等待,也不会引入竞态条件,完全基于ArrayBlockingQueue的阻塞特性实现。

修改后的完整代码

import java.util.concurrent.ArrayBlockingQueue;

public class ParallelExecutor {
    public static void executeInParallel(Runnable[] runnables, int K) {
        if (K <= 0)
            throw new IllegalArgumentException("Can't execute zero or less concurrent runnables");
        if (runnables == null || runnables.length == 0)
            throw new IllegalArgumentException("Must provide something to run");

        ArrayBlockingQueue<Runnable> taskQueue = new ArrayBlockingQueue<>(K);

        // 启动生产者线程
        Thread producer = new Thread(() -> {
            try {
                for (Runnable r : runnables) {
                    taskQueue.put(r);
                }
                // 放入K个null作为结束标记,每个消费者对应一个
                for (int i = 0; i < K; i++) {
                    taskQueue.put(null);
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        });
        producer.start();

        // 启动消费者线程
        K = Math.min(runnables.length, K);
        Thread[] consumers = new Thread[K];
        for (int i = 0; i < K; i++) {
            consumers[i] = new Thread(() -> {
                try {
                    Runnable r;
                    while ((r = taskQueue.take()) != null) {
                        r.run();
                    }
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                } catch (Exception e) {
                    e.printStackTrace();
                }
            });
            consumers[i].start();
        }

        // 等待生产者完成
        try {
            producer.join();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }

        // 等待所有消费者完成
        for (Thread consumer : consumers) {
            try {
                consumer.join();
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
    }

    public static void main(String[] args) {
        Runnable[] tasks = new Runnable[20];
        for (int i = 0; i < tasks.length; i++) {
            final int index = i;
            tasks[i] = () -> {
                System.out.println("Task " + index + " is running.");
                try {
                    Thread.sleep(1000);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
                System.out.println("Task " + index + " is finished.");
            };
        }
        executeInParallel(tasks, 3);
    }
}

关键修改点

  • 生产者线程在提交完所有任务后,新增循环向队列中放入K个null,确保每个消费者线程都能收到结束信号
  • 消费者线程的循环逻辑保持不变,拿到null时自动退出循环,避免无限阻塞

内容的提问来源于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:00