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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:03:25