基于Project Reactor的异步文件逐行读取优化方案咨询
我正在参与Advent of Code活动,想用Project Reactor实现响应式解题方案。目前的代码能部分运行,但遇到缓冲区不足时会出现行被截断的问题。
现有实现代码:
public static Flux<String> getLocalInputForStackOverflow(String filePath) throws IOException { Path dayPath = Path.of(filePath); FileOutputStream resultDay = new FileOutputStream(basePath.resolve("result_day.txt").toFile()); return DataBufferUtils .readAsynchronousFileChannel( () -> AsynchronousFileChannel.open(dayPath), new DefaultDataBufferFactory(), 64) .map(DataBuffer::asInputStream) .map(db -> { try { resultDay.write(db.readAllBytes()); resultDay.write("\n".getBytes()); return db; } catch (FileNotFoundException e) { throw new RuntimeException(e); } catch (IOException e) { throw new RuntimeException(e); } }) .map(InputStreamReader::new) .map(is ->new BufferedReader(is).lines()) .flatMap(Flux::fromStream); }
这段代码原本想以响应式方式逐行读取文件,我通过FileOutputStream把读取内容写入另一个文件对比时,发现缓冲区不足时部分行被截断(try-catch逻辑可忽略)。
已尝试的临时方案:
- 增大缓冲区一次性读取整个文件——非最优解;
- 使用如下代码,但触发警告:
Possibly blocking call in non-blocking context could lead to thread starvation
public static Flux<String> getLocalInput1(int day ) throws IOException { Path dayPath = getFilePath(day); return Flux.using(() -> Files.lines(dayPath), Flux::fromStream, BaseStream::close); }
咨询两个问题:
- 是否存在更优的响应式异步文件读取方式?
- 在使用有限缓冲区的情况下,如何更优地异步逐行读取文件并保证仅读取完整行?
问题1:更优的响应式异步文件读取方式
在Project Reactor生态里,Spring Core的DataBufferUtils配合异步文件通道是标准的非阻塞文件读取方式,但你的实现没处理跨缓冲区的行拼接问题,才导致截断。另外,要避免在map里做阻塞IO操作(比如你的FileOutputStream.write),应该用响应式逻辑处理写入,比如doOnNext或flatMap配合DataBufferUtils.write。
更简洁的非阻塞读取框架示例:
public static Flux<String> readFileReactively(Path filePath) { return DataBufferUtils.read( filePath, new DefaultDataBufferFactory(), 4096) // 推荐用4KB这类合理的缓冲区大小 .transform(Flux::collectList) .flatMapMany(dataBuffers -> { InputStream inputStream = new SequenceInputStream( dataBuffers.stream() .map(DataBuffer::asInputStream) .enumeration()); BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream)); return Flux.fromStream(reader.lines()) .doFinally(signal -> { try { reader.close(); } catch (IOException e) { throw new RuntimeException(e); } }); }); }
注意:如果是超大文件,collectList会占用过多内存,这时候就需要处理行拼接逻辑(对应问题2)。
问题2:有限缓冲区下保证读取完整行
你的核心问题是:异步读取的DataBuffer可能只包含行的一部分,直接转BufferedReader会把截断内容当成完整行。解决思路是维护全局缓冲区缓存未完成的行片段,直到读取到换行符再输出完整行。
具体实现代码:
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferUtils; import org.springframework.core.io.buffer.DefaultDataBufferFactory; import java.io.IOException; import java.nio.charset.StandardCharsets; import java.nio.file.Path; import java.util.ArrayList; import java.util.List; import java.util.concurrent.atomic.AtomicReference; public static Flux<String> readLinesReactively(Path filePath) { AtomicReference<StringBuilder> lineBuffer = new AtomicReference<>(new StringBuilder()); return DataBufferUtils.read( filePath, new DefaultDataBufferFactory(), 64) // 用你指定的小缓冲区测试 .map(dataBuffer -> { byte[] bytes = new byte[dataBuffer.readableByteCount()]; dataBuffer.read(bytes); DataBufferUtils.release(dataBuffer); // 必须释放DataBuffer避免内存泄漏 return new String(bytes, StandardCharsets.UTF_8); }) .flatMap(chunk -> { StringBuilder buffer = lineBuffer.get(); buffer.append(chunk); List<String> lines = new ArrayList<>(); int newlineIndex; while ((newlineIndex = buffer.indexOf("\n")) != -1) { String line = buffer.substring(0, newlineIndex).trim(); // 根据需求决定是否trim lines.add(line); buffer.delete(0, newlineIndex + 1); } lineBuffer.set(buffer); return Flux.fromIterable(lines); }) // 处理最后一行(可能无换行符) .concatWith(Mono.defer(() -> { StringBuilder buffer = lineBuffer.get(); if (buffer.length() > 0) { String lastLine = buffer.toString().trim(); return lastLine.isEmpty() ? Mono.empty() : Mono.just(lastLine); } return Mono.empty(); })); }
这段代码的关键:
- 手动释放
DataBuffer,避免内存泄漏; - 用
AtomicReference维护行缓冲区,保证线程安全; - 逐段分割字符串,提取完整行,缓存不完整片段;
- 流结束时处理无换行符的最后一行。
另外,关于你尝试的Files.lines方案:Files.lines是阻塞IO,在Reactor非阻塞线程池运行会触发警告。如果要使用,需把阻塞操作放到专门的阻塞线程池:
public static Flux<String> getLocalInput1(int day) throws IOException { Path dayPath = getFilePath(day); return Flux.using( () -> Files.lines(dayPath), stream -> Flux.fromStream(stream).publishOn(Schedulers.boundedElastic()), // 切换到阻塞线程池 BaseStream::close ); }
Schedulers.boundedElastic()是Reactor专为阻塞操作设计的线程池,会动态创建线程避免饥饿问题。
内容的提问来源于stack exchange,提问作者robertsci

