Spring Batch的ItemReader如何连接Azure Blob Storage读取容器内所有文件
Spring Batch 读取Azure Blob Storage文件实现方案
1. 引入核心依赖
首先在pom.xml中引入Spring Batch和Azure Storage SDK依赖:
<!-- Spring Batch 依赖 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-batch</artifactId> </dependency> <!-- Azure Blob Storage SDK --> <dependency> <groupId>com.azure</groupId> <artifactId>azure-storage-blob</artifactId> <version>12.25.1</version> </dependency>
2. 配置Azure Blob连接
首先在application.yml中配置Azure Blob的连接信息:
azure: storage: connection-string: 你的Azure存储连接字符串 container-name: 要读取的Blob容器名
然后编写配置类初始化BlobContainerClient:
@Configuration public class AzureBlobConfig { @Value("${azure.storage.connection-string}") private String connectionString; @Value("${azure.storage.container-name}") private String containerName; @Bean public BlobContainerClient blobContainerClient() { return new BlobServiceClientBuilder() .connectionString(connectionString) .buildClient() .getBlobContainerClient(containerName); } }
3. 自定义Azure Blob ItemReader
我们需要自定义ItemReader实现类,批量拉取容器内的Blob文件,逐个返回给后续Processor处理:
@StepScope @Component public class AzureBlobItemReader implements ItemReader<BlobClient> { private final BlobContainerClient blobContainerClient; private Iterator<BlobItem> blobItemIterator; // 构造注入Blob容器客户端 public AzureBlobItemReader(BlobContainerClient blobContainerClient) { this.blobContainerClient = blobContainerClient; // 首次调用时加载容器内所有Blob列表(可添加prefix参数过滤指定前缀的文件) this.blobItemIterator = blobContainerClient.listBlobs().iterator(); } @Override public BlobClient read() { // 遍历到末尾返回null,标识Reader读取结束 if (!blobItemIterator.hasNext()) { return null; } BlobItem currentBlob = blobItemIterator.next(); // 跳过文件夹(如果不需要可以去掉该判断) if (currentBlob.isPrefix()) { return read(); } // 返回BlobClient对象,可直接获取文件流、元数据等信息 return blobContainerClient.getBlobClient(currentBlob.getName()); } }
如果需要直接返回文件内容而非BlobClient,可在read方法中调用downloadContent()方法读取文件内容后返回自定义的业务对象即可
4. 配置Spring Batch Step与Job
将自定义的Reader绑定到Step流程中,和你自己的Processor、Writer组合使用:
@Configuration @EnableBatchProcessing public class BatchJobConfig { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; private final AzureBlobItemReader azureBlobItemReader; // 你自己实现的ItemProcessor private final YourCustomProcessor customProcessor; // 你自己实现的ItemWriter private final YourCustomWriter customWriter; // 构造注入所有Bean public BatchJobConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, AzureBlobItemReader azureBlobItemReader, YourCustomProcessor customProcessor, YourCustomWriter customWriter) { this.jobBuilderFactory = jobBuilderFactory; this.stepBuilderFactory = stepBuilderFactory; this.azureBlobItemReader = azureBlobItemReader; this.customProcessor = customProcessor; this.customWriter = customWriter; } @Bean public Step processBlobStep() { return stepBuilderFactory.get("processBlobStep") // 可根据需要调整chunk大小 .<BlobClient, YourBusinessEntity>chunk(10) .reader(azureBlobItemReader) .processor(customProcessor) .writer(customWriter) .build(); } @Bean public Job blobProcessJob() { return jobBuilderFactory.get("blobProcessJob") .start(processBlobStep()) .build(); } }
注意事项
- 处理大体积Blob文件时,建议使用分片下载接口而非一次性读取全量内容,避免内存溢出
- 可在
listBlobs()方法中传入ListBlobsOptions参数,实现按文件前缀、修改时间等规则过滤需要处理的文件 - 生产环境建议添加异常处理机制,捕获Blob读取异常后根据业务需求选择重试或者跳过异常文件
内容的提问来源于stack exchange,提问作者Krish
相关产品推荐
相关产品推荐

