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

使用CompletableFuture处理readLine阻塞:超时后如何获取部分结果?

Handling Partial Results from Blocking 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 check Thread.currentThread().isInterrupted() inside the loop to exit early instead of continuing to block on readLine().
  • Stream Handling: Using try-with-resources ensures the BufferedReader is closed properly, even if reading is interrupted. If your InputStream is 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 the InputStream in 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, ConcurrentLinkedQueue will 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:09:35