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

Spring Batch中InputStreamResource读取报错及动态读取文件问题求助

Spring Batch FlatFileItemReader 使用 InputStreamResource 报错及动态读取问题

问题背景

我在Spring Batch作业中定义了以下FlatFileItemReader:

@Bean
public FlatFileItemReader<BookingInfo> bookingInfoReader() {
  FlatFileItemReader<BookingInfo> itemReader = new FlatFileItemReader<>();
  ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
  BlobClient blobClient =
      blobContainerClient.getBlobClient("path/file.csv");
  blobClient.downloadStream(outputStream);
  final byte[] bytes = outputStream.toByteArray();
  ByteArrayInputStream inputStream = new ByteArrayInputStream(bytes);
  InputStreamResource resource = new InputStreamResource(inputStream);
  itemReader.setResource(resource);
  itemReader.setName("csvReader");
  itemReader.setLinesToSkip(1);
  itemReader.setLineMapper(lineMapper());
  return itemReader;
}

运行时出现错误:

java.lang.IllegalStateException: InputStream has already been read - do not use InputStreamResource if a stream needs to be read multiple times

改成ByteArrayResource后错误消失,但出现新问题:更新文件内容后,作业始终读取旧内容,只有重启Pod才会获取新内容:

@Bean
public FlatFileItemReader<BookingInfo> bookingInfoReader() {

  FlatFileItemReader<BookingInfo> itemReader = new FlatFileItemReader<>();
  ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
  BlobClient blobClient =
      blobContainerClient.getBlobClient("path/file.csv");
  blobClient.downloadStream(outputStream);
  final byte[] bytes = outputStream.toByteArray();
  ByteArrayInputStream inputStream = new ByteArrayInputStream(bytes);
  
  ByteArrayResource byteArrayResource = new ByteArrayResource(inputStream.readAllBytes());
  itemReader.setResource(byteArrayResource); 

  itemReader.setName("csvReader");
  itemReader.setLinesToSkip(1);
  itemReader.setLineMapper(lineMapper());
  return itemReader;
}

我希望回到InputStreamResource的实现并修复报错,同时解决动态读取最新文件内容的问题。

问题原因

  1. InputStreamResource 报错原因:它包装的InputStream是一次性的,只能被读取一次。Spring Batch在作业执行过程中(比如作业重启、内部校验逻辑)可能会多次尝试读取资源流,流被耗尽后就会抛出异常。
  2. ByteArrayResource 静态读取问题:Bean是单例模式,初始化时就完成了Blob文件下载,后续作业执行只会复用这部分已加载的旧内容,无法感知文件更新。

解决方案

1. 自定义支持动态读取的Resource

实现一个自定义Resource,每次调用getInputStream()时都重新从Blob存储下载最新文件并返回新的流:

public class BlobResource extends AbstractResource {

    private final BlobContainerClient blobContainerClient;
    private final String blobPath;

    public BlobResource(BlobContainerClient blobContainerClient, String blobPath) {
        this.blobContainerClient = blobContainerClient;
        this.blobPath = blobPath;
    }

    @Override
    public InputStream getInputStream() throws IOException {
        ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
        BlobClient blobClient = blobContainerClient.getBlobClient(blobPath);
        blobClient.downloadStream(outputStream);
        return new ByteArrayInputStream(outputStream.toByteArray());
    }

    @Override
    public String getDescription() {
        return "Blob resource located at path: " + blobPath;
    }
}

2. 修改Reader Bean配置

将Reader的资源替换为自定义的BlobResource,这样每次需要读取流时都会重新下载最新文件,同时每次返回新的InputStream实例:

@Bean
public FlatFileItemReader<BookingInfo> bookingInfoReader() {
    FlatFileItemReader<BookingInfo> itemReader = new FlatFileItemReader<>();
    // 使用自定义BlobResource,延迟文件下载到实际读取时
    BlobResource blobResource = new BlobResource(blobContainerClient, "path/file.csv");
    itemReader.setResource(blobResource);
    itemReader.setName("csvReader");
    itemReader.setLinesToSkip(1);
    itemReader.setLineMapper(lineMapper());
    return itemReader;
}

方案说明

  • 自定义BlobResource解决了两个核心问题:每次读取都获取最新文件内容,且每次返回新的InputStream实例,支持多次读取。
  • 把文件下载逻辑从Bean初始化阶段延迟到实际读取流时,避免了单例Bean导致的静态缓存问题。

内容的提问来源于stack exchange,提问作者Rahul Raj

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 15:53:34