为进程中的阻塞式read()添加超时中断的实现问题
问题分析
你的代码里,output.readLine()是阻塞式调用——如果输入流没有数据,这个方法会一直挂着不返回。而超时判断的代码是在readLine()执行完之后才运行的,所以只要readLine()没结束,超时逻辑根本触发不了,这就是问题的核心。
解决方案
要给阻塞式的read操作加超时,核心是让阻塞的读取操作能被超时中断。下面提供两种可行的实现方式:
方法1:线程+Future实现可超时读取
把读取逻辑放到子线程中,用Future.get(timeout)控制超时,超时后直接取消任务、中断线程并关闭输入流,强制终止read操作。
import java.io.BufferedReader; import java.io.InputStream; import java.io.InputStreamReader; import java.util.ArrayList; import java.util.List; import java.util.Objects; import java.util.concurrent.*; import java.util.stream.Collectors; public class TimeoutReader { public List<String> readWithTimeout(InputStream inputStream, Long timeout, TimeUnit timeUnit) throws Exception { if (Objects.isNull(timeUnit) || Objects.isNull(timeout)) { try (BufferedReader output = new BufferedReader(new InputStreamReader(inputStream))) { return output.lines().collect(Collectors.toList()); } } // 创建单线程池执行读取任务 ExecutorService executor = Executors.newSingleThreadExecutor(); Future<List<String>> future = executor.submit(() -> { try (BufferedReader output = new BufferedReader(new InputStreamReader(inputStream))) { List<String> result = new ArrayList<>(); String line; while ((line = output.readLine()) != null) { result.add(line); } return result; } }); try { // 等待任务完成,超时则抛出异常 return future.get(timeout, timeUnit); } catch (TimeoutException e) { // 超时后取消任务,中断线程 future.cancel(true); // 关闭输入流强制终止read操作 inputStream.close(); throw new ShellTaskExecutionException(String.format("Task reached timeout %s %s", timeout, timeUnit)); } finally { executor.shutdownNow(); } } }
方法2:NIO可中断输入流(适配支持中断的场景)
将传统InputStream转为NIO的InterruptibleChannel包装流,这样线程被中断时,read操作会抛出ClosedByInterruptException,从而终止阻塞。
import java.io.BufferedReader; import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; import java.nio.channels.Channels; import java.nio.channels.ReadableByteChannel; import java.util.ArrayList; import java.util.List; import java.util.Objects; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; public class NioTimeoutReader { public List<String> readWithTimeout(InputStream inputStream, Long timeout, TimeUnit timeUnit) throws Exception { if (Objects.isNull(timeUnit) || Objects.isNull(timeout)) { try (BufferedReader output = new BufferedReader(new InputStreamReader(inputStream))) { return output.lines().collect(Collectors.toList()); } } ReadableByteChannel channel = Channels.newChannel(inputStream); InputStream nioInputStream = Channels.newInputStream(channel); try (BufferedReader output = new BufferedReader(new InputStreamReader(nioInputStream))) { List<String> result = new ArrayList<>(); long timeoutMillis = timeUnit.toMillis(timeout); long startTime = System.currentTimeMillis(); Thread currentThread = Thread.currentThread(); // 启动超时中断线程 Thread timeoutThread = new Thread(() -> { try { Thread.sleep(timeoutMillis); currentThread.interrupt(); channel.close(); } catch (InterruptedException | IOException e) { // 线程被中断或通道已关闭,无需额外处理 } }); timeoutThread.start(); try { String line; while ((line = output.readLine()) != null) { result.add(line); // 重置超时时间(如果需要累计超时则注释此行) startTime = System.currentTimeMillis(); timeoutThread.interrupt(); timeoutThread = new Thread(() -> { try { Thread.sleep(timeoutMillis - (System.currentTimeMillis() - startTime)); currentThread.interrupt(); channel.close(); } catch (InterruptedException | IOException e) { } }); timeoutThread.start(); } // 读取完成,中断超时线程 timeoutThread.interrupt(); return result; } catch (IOException e) { // 判断是否为超时中断导致的异常 if (Thread.currentThread().isInterrupted()) { throw new ShellTaskExecutionException(String.format("Task reached timeout %s %s", timeout, timeUnit)); } throw e; } } } }
注意事项
- 方法1是通用方案,适配所有输入流场景,推荐优先使用。
- 如果读取的是子进程输出流,超时后建议额外调用
Process.destroy()销毁子进程,避免资源泄漏。 - 两种实现都要确保输入流被正确关闭,防止资源泄漏。
内容的提问来源于stack exchange,提问作者itzik
相关产品推荐
相关产品推荐

