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的实现并修复报错,同时解决动态读取最新文件内容的问题。
问题原因
- InputStreamResource 报错原因:它包装的InputStream是一次性的,只能被读取一次。Spring Batch在作业执行过程中(比如作业重启、内部校验逻辑)可能会多次尝试读取资源流,流被耗尽后就会抛出异常。
- 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
相关产品推荐
相关产品推荐

