如何使用Spring Batch将文件分块处理并写入多个文件
Spring Batch 按输入文件组分块输出多CSV实现方案
核心逻辑不需要调整现有MultiResourceItemReader的读取配置,也不需要提前把50个CSV手动分组,全程保持流式读写,不会产生全量数据加载的内存压力,完全适配6000万条数据的处理量级。
实现思路
放弃原来绑定单个输出资源的FlatFileItemWriter,替换为可按输入文件处理进度动态切换输出目标的多资源写出器,配合读监听器统计已处理的输入文件个数,每累计处理完3个输入文件就自动切换到新的分块CSV写入即可。
具体实现步骤
1. 注册Step级上下文计数变量
在Step执行上下文中维护两个计数变量,用于支撑切分逻辑,同时支持作业断点续跑:
processedInputFileCount:已完成读取的输入文件个数,初始值为0currentOutputChunkIndex:当前输出的分块文件序号,初始值为1
2. 自定义按输入文件数触发切分的多资源写出器
Spring Batch原生MultiResourceItemWriter默认按写入数据条数切分文件,只需要重写它的资源切换判断逻辑,就能适配按输入文件个数切分的需求,不需要从零实现写出逻辑:
public class FileGroupBoundedWriter extends MultiResourceItemWriter<DataRow> { private final int FILE_GROUP_SIZE = 3; private StepExecution stepExecution; @BeforeStep public void bindStepContext(StepExecution stepExecution) { this.stepExecution = stepExecution; } @Override protected Resource getNewOutputResource() { int processedFileCount = stepExecution.getExecutionContext() .getInt("processedInputFileCount", 0); // 已处理文件数达到3的倍数时,切换到新的输出文件 if (processedFileCount > 0 && processedFileCount % FILE_GROUP_SIZE == 0) { int currentChunkIndex = stepExecution.getExecutionContext() .getInt("currentOutputChunkIndex", 1); currentChunkIndex++; stepExecution.getExecutionContext() .putInt("currentOutputChunkIndex", currentChunkIndex); return new FileSystemResource( String.format("merged-chunk-%d.csv", currentChunkIndex) ); } // 未达阈值时继续写入当前文件 return null; } }
该写出器内部委派的
FlatFileItemWriter直接复用你之前单文件输出时的行映射、字段格式化配置,不需要调整输出格式逻辑。
3. 添加读监听器同步输入文件处理进度
实现ItemReadListener,在每次读操作后判断是否切换到了新的输入文件,实时更新上下文中的已处理文件计数:
public class InputFileProgressListener implements ItemReadListener<DataRow> { private StepExecution stepExecution; private Resource lastReadResource; private final MultiResourceItemReader<DataRow> multiReader; // 构造注入已经配置好的多文件读取器 public InputFileProgressListener(MultiResourceItemReader<DataRow> multiReader) { this.multiReader = multiReader; } @BeforeStep public void initContext(StepExecution stepExecution) { this.stepExecution = stepExecution; stepExecution.getExecutionContext().putInt("processedInputFileCount", 0); stepExecution.getExecutionContext().putInt("currentOutputChunkIndex", 1); } @Override public void afterRead(DataRow item) { Resource currentResource = multiReader.getCurrentResource(); // 检测到输入文件切换时,已处理文件数+1 if (lastReadResource != null && !currentResource.equals(lastReadResource)) { int count = stepExecution.getExecutionContext().getInt("processedInputFileCount", 0); stepExecution.getExecutionContext().putInt("processedInputFileCount", count + 1); } lastReadResource = currentResource; } // 其余监听器方法留空即可 @Override public void beforeRead() {} @Override public void onReadError(Exception ex) {} }
边界处理说明
- 最后一组如果剩余输入文件不足3个,会自动写入最后一个分块CSV,不需要额外编写边界判断逻辑
- 所有计数变量存在Step执行上下文,作业中断重启时会自动恢复进度,不会出现文件重复合并、数据遗漏的问题
- 整个流程保持逐行读写的流式处理逻辑,内存开销和之前单文件合并的方案基本一致,不会因为数据量达到6000万出现OOM
内容的提问来源于stack exchange,提问作者Prateek Singh
相关产品推荐
相关产品推荐

