bufferedwriter 不能直接用于 reactive stream,因其阻塞调用、不支持背压且缺乏响应式语义;应改用 databuffer + filechannel 或 files.write(publisher) 实现非阻塞缓存写入,并通过 buffer() 和 subscribeon(boundedelastic) 控制缓存与线程隔离。

为什么不能直接用 BufferedWriter 包装 Flux
• BufferedWriter.write() 和 flush() 都是阻塞调用,在 Reactor 的 event loop(如 elastic 或 parallel scheduler)中调用会阻塞线程,破坏响应式流的非阻塞特性。
• 它不支持背压:无法根据下游消费速度控制上游生产节奏,容易 OOM。
• 没有 onErrorResume / onBackpressureBuffer 等响应式语义,异常传播和缓冲策略需手动模拟,极易出错。
推荐方案:用 DataBuffer + FileChannel + Mono/Flux 实现非阻塞缓存写入
Spring WebFlux 和 Project Reactor 提供了真正的响应式文件 I/O 支持:
• 使用 Files.write(Path, Publisher
• 或通过 AsynchronousFileChannel + Mono.fromCompletionStage 封装写入操作
• 缓存逻辑应由上游控制(如用 flux.buffer(1024).map(lines → String.join("\n", lines)) 聚合行),而非依赖 BufferedWriter 的内部缓冲
实际可落地的缓存写入示例(Project Reactor)
假设你要把字符串流按 100 行一批写入文件:
Flux<string> lines = Flux.range(1, 10000)
.map(i -> "log-" + i)
.delayElements(Duration.ofMillis(1)); // 模拟异步数据源
Path path = Paths.get("output.log");
// 缓存 100 行 → 合并为一个字符串 → 写入文件(追加)
lines.buffer(100)
.map(batch -> String.join("\n", batch) + "\n")
.flatMap(content -> Mono.fromCallable(() -> {
Files.write(path, content.getBytes(StandardCharsets.UTF_8),
StandardOpenOption.CREATE, StandardOpenOption.APPEND);
return content;
}).subscribeOn(Schedulers.boundedElastic())) // 必须切换到 IO 线程池
.blockLast(); // 仅演示;生产环境用 subscribe()
</string>
关键点:
• buffer(100) 实现应用层缓存,可控且支持背压
• subscribeOn(Schedulers.boundedElastic()) 隔离阻塞 I/O,避免污染主线程
• 不用 BufferedWriter,而是用 Files.write(... APPEND) 或 AsynchronousFileChannel 做真正异步写入
如果必须用传统 FileWriter 场景(如遗留系统集成)
只能退回到“响应式外壳 + 阻塞内核”模式,但要明确代价:
• 用 Mono.fromCallable 封装 BufferedWriter 写入逻辑
• 强制调度到 Schedulers.boundedElastic()
• 手动处理背压(如用 onBackpressureBuffer(1000) 限制缓存队列大小)
• 注意 close() 必须在 finally 块中确保执行,避免资源泄漏
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











