Java:单生产者多消费者硬件模拟任务的线程控制问题
解决单生产者-多消费者场景下的Future阻塞提交问题
你的核心痛点很明确:主线程如果通过while(!future.isDone())轮询等待任务完成,会导致后续设备的任务无法及时提交,完全丧失了多设备并行执行的优势,最终延迟不断累积。我们可以通过异步回调替代主动轮询的方式解决这个问题,不用阻塞主线程,也不需要额外新建轮询线程。
方案1:使用CompletableFuture(Java 8+推荐)
CompletableFuture是Java 8引入的增强版Future,支持链式异步回调,任务完成时会自动触发后续逻辑,完全不需要主线程主动轮询。
代码示例:
// 假设每个设备有自己的SingleThreadExecutorService ExecutorService device1Executor = Executors.newSingleThreadExecutor(); ExecutorService device2Executor = Executors.newSingleThreadExecutor(); // 提交设备1的任务,异步处理完成事件 CompletableFuture<Void> future1 = CompletableFuture.runAsync(() -> { // 模拟硬件任务,比如sleep 1000ms try { Thread.sleep(1000); System.out.println("设备1任务完成"); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }, device1Executor); // 立即提交设备2的任务,无需等待设备1完成 CompletableFuture<Void> future2 = CompletableFuture.runAsync(() -> { try { Thread.sleep(800); System.out.println("设备2任务完成"); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }, device2Executor); // 如果主线程需要等待所有任务完成再继续,用allOf(此时所有任务已提交,阻塞不会影响并行) CompletableFuture.allOf(future1, future2).join(); // 最后记得关闭线程池(根据你的业务生命周期调整) device1Executor.shutdown(); device2Executor.shutdown();
优势:
- 主线程可以一次性提交所有设备的任务,真正实现多设备并行执行
- 任务完成后的逻辑(比如结果处理、状态更新)可以直接通过
thenAccept/thenRun等方法异步触发 - 支持异常处理(
exceptionally方法),比普通Future更健壮
方案2:使用FutureTask+监听器(兼容Java 7及以下)
如果你的项目无法升级到Java 8,可以用FutureTask的addListener方法,给任务添加完成监听器,任务结束时自动执行回调逻辑。
代码示例:
ExecutorService device1Executor = Executors.newSingleThreadExecutor(); ExecutorService device2Executor = Executors.newSingleThreadExecutor(); // 可以用一个专门的线程池处理回调,避免占用设备线程或主线程 ExecutorService callbackExecutor = Executors.newCachedThreadPool(); // 设备1的任务 Runnable task1 = () -> { try { Thread.sleep(1000); System.out.println("设备1任务完成"); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }; FutureTask<Void> futureTask1 = new FutureTask<>(task1, null); // 添加任务完成监听器 futureTask1.addListener(() -> { // 任务完成后的处理逻辑 System.out.println("设备1任务回调触发"); }, callbackExecutor); device1Executor.submit(futureTask1); // 立即提交设备2的任务 Runnable task2 = () -> { try { Thread.sleep(800); System.out.println("设备2任务完成"); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }; FutureTask<Void> futureTask2 = new FutureTask<>(task2, null); futureTask2.addListener(() -> { System.out.println("设备2任务回调触发"); }, callbackExecutor); device2Executor.submit(futureTask2); // 等待所有任务完成(可选) try { futureTask1.get(); futureTask2.get(); } catch (InterruptedException | ExecutionException e) { e.printStackTrace(); } // 关闭线程池 device1Executor.shutdown(); device2Executor.shutdown(); callbackExecutor.shutdown();
优势:
- 兼容低版本Java
- 同样避免了主线程阻塞轮询,任务提交完全并行
- 可以自定义回调线程池,灵活控制回调的执行资源
核心思路总结
你之前的问题本质是把任务提交和结果等待串行化了,导致主线程被阻塞在等待上,无法及时提交后续任务。改用异步回调后,主线程只负责"提交任务"这一件事,任务完成的通知由JDK异步触发,彻底解决了提交延迟的问题,同时保留了每个设备任务串行、多设备任务并行的设计优势。
内容的提问来源于stack exchange,提问作者Bontempos
相关产品推荐
相关产品推荐

