Spring WebFlux/Project Reactor:如何将Flux数据写入文本文件?
问题根源
你代码的问题出在两个核心点:
- 异步时序错误:
subscribe是异步触发的,调用后代码会立刻执行writer.flush()和writer.close()——此时Flux还没开始发射数据,文件已经被关闭,自然写不进内容。 - 阻塞IO破坏非阻塞特性:
PrintWriter是阻塞式IO工具,在Reactor的异步线程中调用它的方法,会把非阻塞流程强行变成阻塞,完全违背WebFlux的设计原则。
非阻塞写入文件的最佳方案
推荐使用Spring WebFlux生态中的DataBufferUtils工具,它专门适配非阻塞IO场景,能完美对接Flux的异步流处理:
import org.springframework.core.io.FileSystemResource; import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DefaultDataBufferFactory; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.core.io.buffer.DataBufferUtils; import java.nio.charset.StandardCharsets; import java.nio.file.Paths; import java.nio.file.StandardOpenOption; // 省略Students类定义 public class NonBlockingFileWriter { public static void main(String[] args) { WebClient webClient = WebClient.create("你的服务基础URL"); FileSystemResource targetFile = new FileSystemResource(Paths.get("ex1.txt")); webClient.get().uri("/users/listall") .retrieve() .bodyToFlux(Students.class) .map(stud -> stud.getName() + ", " + stud.getAge() + "\n") // 拼接内容并添加换行 .map(content -> { byte[] bytes = content.getBytes(StandardCharsets.UTF_8); return DefaultDataBufferFactory.sharedInstance.wrap(bytes); }) .flatMap(dataBuffer -> DataBufferUtils.write(dataBuffer, targetFile, StandardOpenOption.CREATE)) .doOnComplete(() -> System.out.println("所有数据已异步写入文件")) .subscribe(); } }
方案说明
DataBufferUtils.write是纯非阻塞操作,会将DataBuffer内容异步写入文件,不会阻塞Reactor的事件循环线程- 无需手动关闭文件资源,
DataBufferUtils会自动处理资源释放 - 全程保持WebFlux的非阻塞特性,完全契合异步流的设计逻辑
原生NIO替代方案(无Spring依赖)
如果不想依赖Spring的DataBuffer体系,可以用JDK原生的AsynchronousFileChannel实现异步写入:
import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.nio.file.Path; import java.nio.file.Paths; import java.nio.channels.AsynchronousFileChannel; import java.nio.file.StandardOpenOption; import java.util.concurrent.atomic.AtomicLong; import reactor.core.publisher.Mono; import org.springframework.web.reactive.function.client.WebClient; // 省略Students类定义 public class JdkAsyncFileWriter { public static void main(String[] args) throws Exception { WebClient webClient = WebClient.create("你的服务基础URL"); Path filePath = Paths.get("ex1.txt"); AsynchronousFileChannel fileChannel = AsynchronousFileChannel.open( filePath, StandardOpenOption.CREATE, StandardOpenOption.WRITE ); AtomicLong writePosition = new AtomicLong(0); webClient.get().uri("/users/listall") .retrieve() .bodyToFlux(Students.class) .map(stud -> stud.getName() + ", " + stud.getAge() + "\n") .map(content -> content.getBytes(StandardCharsets.UTF_8)) .flatMap(bytes -> { long pos = writePosition.getAndAdd(bytes.length); return Mono.fromFuture(fileChannel.write(ByteBuffer.wrap(bytes), pos)); }) .doOnComplete(() -> { try { fileChannel.close(); System.out.println("文件异步写入完成"); } catch (Exception e) { e.printStackTrace(); } }) .subscribe(); } }
方案说明
AsynchronousFileChannel是JDK原生的异步文件通道,所有写入操作均为非阻塞- 通过
Mono.fromFuture将异步通道的Future结果转换为Reactor的Mono,保持流的异步特性 - 用
AtomicLong记录写入位置,避免多线程环境下的位置错乱
内容的提问来源于stack exchange,提问作者ruiftmota
相关产品推荐
相关产品推荐

