使用CompletableFuture处理readLine阻塞:超时后如何获取部分结果?
readLine() with CompletableFuture Timeouts Great question—this is a common pain point when dealing with blocking I/O and async timeouts, since standard CompletableFuture cancellation doesn’t preserve partial results out of the box. The core solution is to decouple the reading operation from result storage: use a thread-safe container to capture lines as they’re read, then access that container even if the future times out.
Step-by-Step Implementation
1. Create a Thread-Safe Partial Result Holder
First, we need a container to store lines as they’re read. This needs to be thread-safe because the async reading thread will write to it, and the main thread will read from it if a timeout occurs.
import java.util.Queue; import java.util.concurrent.ConcurrentLinkedQueue; public class PartialLineReaderResult { private final Queue<String> capturedLines = new ConcurrentLinkedQueue<>(); private volatile boolean readingCompleted = false; // Add a newly read line to the container public void captureLine(String line) { capturedLines.add(line); } // Get all lines read so far, joined into a single string public String getPartialContent() { return String.join(System.lineSeparator(), capturedLines); } public void markReadingCompleted() { readingCompleted = true; } public boolean isReadingCompleted() { return readingCompleted; } }
We use ConcurrentLinkedQueue here because it’s thread-safe and avoids explicit synchronization overhead for individual adds.
2. Wrap the Blocking Read with Async Logic
Next, we’ll wrap the readLine() loop in a CompletableFuture, updating our result holder with every line. We also add checks to stop reading if the thread is interrupted (triggered when we cancel the future on timeout).
import java.io.BufferedReader; import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; public class TimeoutAwareLineReader { public String readLinesWithTimeout(InputStream inputStream, long timeout, TimeUnit timeUnit) throws InterruptedException, TimeoutException { PartialLineReaderResult resultHolder = new PartialLineReaderResult(); ExecutorService executor = Executors.newSingleThreadExecutor(); // Single thread for safe sequential reading CompletableFuture<String> readFuture = CompletableFuture.supplyAsync(() -> { try (BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream))) { String line; while ((line = reader.readLine()) != null) { // Check if we've been interrupted (future was cancelled) if (Thread.currentThread().isInterrupted()) { break; } resultHolder.captureLine(line); } resultHolder.markReadingCompleted(); return resultHolder.getPartialContent(); } catch (IOException e) { // If the exception is caused by an interrupt (e.g., stream closed due to timeout), return partial results if (Thread.currentThread().isInterrupted()) { return resultHolder.getPartialContent(); } // For other IO errors, wrap and rethrow throw new CompletionException(e); } }, executor); try { // Wait for the read to finish or timeout return readFuture.orTimeout(timeout, timeUnit).join(); } catch (TimeoutException e) { // Cancel the future and interrupt the reading thread to stop further I/O readFuture.cancel(true); executor.shutdownNow(); // Return whatever lines we managed to read before the timeout return resultHolder.getPartialContent(); } finally { // Ensure the executor is cleaned up even if no timeout occurs if (!executor.isShutdown()) { executor.shutdown(); } } } }
Key Details to Note
- Thread Interruption: When we call
readFuture.cancel(true), it interrupts the thread running the read loop. We checkThread.currentThread().isInterrupted()inside the loop to exit early instead of continuing to block onreadLine(). - Stream Handling: Using try-with-resources ensures the
BufferedReaderis closed properly, even if reading is interrupted. If yourInputStreamis shared or needs to stay open, adjust the resource management accordingly. - Executor Safety: We use a single-threaded executor to avoid concurrent reads on the same
InputStream, which could cause corrupted data or unexpected behavior.
Edge Cases to Consider
- Slow/Unresponsive Streams: If
readLine()doesn’t respond to thread interrupts (e.g., some network streams), you may need to explicitly close theInputStreamin the timeout block to force the read to exit. Be cautious with this if the stream is used elsewhere. - High-Volume Reads: For extremely fast streams,
ConcurrentLinkedQueuewill handle the load, but if you need to minimize memory usage, consider limiting the number of stored lines (e.g., keep only the last N lines).
内容的提问来源于stack exchange,提问作者Mauricio Villegas

