如何通过FlatFileItemWriter向InputStream写入数据并解决报错?
问题分析与解决方案
这个坑我之前踩过!核心问题在于FlatFileItemWriter从设计之初就是绑定文件系统的——它在open()阶段会调用getResource().getFile()来获取文件对象,而InputStreamResource本质上只是一个流的包装,没有对应的物理文件路径,所以必然会抛出FileNotFoundException。
你的需求是直接通过S3Client把数据以流的形式上传到S3,完全不需要落地文件,所以根本没必要继承FlatFileItemWriter,直接自定义一个适配S3的ItemWriter才是正确的思路。下面给你两种可行的实现方式:
方案一:内存暂存后批量上传(适合中小数据量)
这种方式先把数据写到内存流里,待所有数据处理完成后,一次性上传到S3,实现简单,适合数据量不大的场景:
import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.PutObjectRequest; import org.springframework.batch.item.ItemWriter; import org.springframework.batch.item.ItemStream; import org.springframework.batch.item.ItemStreamException; import org.springframework.batch.item.file.transform.LineAggregator; import org.springframework.batch.item.ExecutionContext; import java.io.*; import java.nio.charset.StandardCharsets; import java.util.List; public class S3BatchItemWriter<T> implements ItemWriter<T>, ItemStream { private final S3Client s3Client; private final String bucketName; private final String s3Key; private final LineAggregator<T> lineAggregator; // 用来把Item转换成一行字符串,复用FlatFile的逻辑 private ByteArrayOutputStream byteOutputStream; private BufferedWriter bufferedWriter; // 构造函数注入所有依赖 public S3BatchItemWriter(S3Client s3Client, String bucketName, String s3Key, LineAggregator<T> lineAggregator) { this.s3Client = s3Client; this.bucketName = bucketName; this.s3Key = s3Key; this.lineAggregator = lineAggregator; } @Override public void open(ExecutionContext executionContext) throws ItemStreamException { // 初始化内存流和写入器 byteOutputStream = new ByteArrayOutputStream(); bufferedWriter = new BufferedWriter(new OutputStreamWriter(byteOutputStream, StandardCharsets.UTF_8)); } @Override public void write(List<? extends T> items) throws Exception { // 遍历Item,转换成字符串写入流 for (T item : items) { String line = lineAggregator.aggregate(item); bufferedWriter.write(line); bufferedWriter.newLine(); } bufferedWriter.flush(); } @Override public void update(ExecutionContext executionContext) throws ItemStreamException { // 可选:保存执行状态,用于任务重启(比如记录已写入的Item数量) } @Override public void close() throws ItemStreamException { try { if (bufferedWriter != null) { bufferedWriter.close(); } // 将内存流转为InputStream,上传到S3 try (InputStream inputStream = new ByteArrayInputStream(byteOutputStream.toByteArray())) { PutObjectRequest request = PutObjectRequest.builder() .bucket(bucketName) .key(s3Key) .build(); s3Client.putObject(request, RequestBody.fromInputStream(inputStream, byteOutputStream.size())); } } catch (IOException e) { throw new ItemStreamException("上传数据到S3失败", e); } } }
优势:
- 完全复用Spring Batch的
LineAggregator逻辑,不用重新写Item转字符串的代码 - 实现简单,没有复杂的线程处理
- 不需要任何本地文件,纯内存操作
方案二:边写边上传(适合大数据量)
如果你的数据量很大,内存流会占用过多内存导致OOM,这时候可以用**管道流(PipedInputStream/PipedOutputStream)**实现边写边上传,数据不会在内存中堆积:
import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.PutObjectRequest; import org.springframework.batch.item.ItemWriter; import org.springframework.batch.item.ItemStream; import org.springframework.batch.item.ItemStreamException; import org.springframework.batch.item.file.transform.LineAggregator; import org.springframework.batch.item.ExecutionContext; import java.io.*; import java.nio.charset.StandardCharsets; import java.util.List; public class S3StreamingItemWriter<T> implements ItemWriter<T>, ItemStream { private final S3Client s3Client; private final String bucketName; private final String s3Key; private final LineAggregator<T> lineAggregator; private BufferedWriter bufferedWriter; public S3StreamingItemWriter(S3Client s3Client, String bucketName, String s3Key, LineAggregator<T> lineAggregator) { this.s3Client = s3Client; this.bucketName = bucketName; this.s3Key = s3Key; this.lineAggregator = lineAggregator; } @Override public void open(ExecutionContext executionContext) throws ItemStreamException { try { // 创建管道流:一边写数据到输出流,另一边从输入流读取上传 PipedInputStream pipedIn = new PipedInputStream(); PipedOutputStream pipedOut = new PipedOutputStream(pipedIn); this.bufferedWriter = new BufferedWriter(new OutputStreamWriter(pipedOut, StandardCharsets.UTF_8)); // 启动单独线程处理S3上传 new Thread(() -> { try { PutObjectRequest request = PutObjectRequest.builder() .bucket(bucketName) .key(s3Key) .build(); // 这里传入-1表示不指定内容长度,S3会自动处理 s3Client.putObject(request, RequestBody.fromInputStream(pipedIn, -1)); } catch (Exception e) { throw new RuntimeException("S3上传线程异常", e); } finally { try { pipedIn.close(); } catch (IOException ignored) {} } }).start(); } catch (IOException e) { throw new ItemStreamException("初始化管道流失败", e); } } @Override public void write(List<? extends T> items) throws Exception { for (T item : items) { String line = lineAggregator.aggregate(item); bufferedWriter.write(line); bufferedWriter.newLine(); } bufferedWriter.flush(); } @Override public void update(ExecutionContext executionContext) throws ItemStreamException { // 可选:保存执行状态 } @Override public void close() throws ItemStreamException { try { if (bufferedWriter != null) { bufferedWriter.close(); // 关闭写入器会触发管道输出流的EOF,上传线程会自动结束 } } catch (IOException e) { throw new ItemStreamException("关闭写入器失败", e); } } }
注意事项:
- 上传线程要单独启动,避免阻塞主线程的Item处理逻辑
- 要处理好线程的异常捕获和资源关闭,防止内存泄漏
- 如果需要支持任务重启,需要额外处理状态保存和断点续传逻辑
为什么不推荐强行适配FlatFileItemWriter?
FlatFileItemWriter的核心逻辑都是围绕文件展开的,比如文件追加、分片写入、重启时的文件位置恢复等,这些功能都依赖物理文件的存在。强行用InputStreamResource去适配,不仅会抛出异常,就算通过自定义Resource绕过getFile()检查,后续的逻辑也会因为没有文件而出现各种问题,完全得不偿失。
内容的提问来源于stack exchange,提问作者bumi25
相关产品推荐
相关产品推荐

