Spring Batch大数据多文件写入优化咨询:现有Chunk写法不够高效
优化Spring Batch多文件Chunk写入方案
针对你当前为每个文件创建独立Step的低效问题,这里提供两种更灵活高效的方案:
方案一:单Step+动态ItemWriter
核心思路是用一个统一的Reader读取所有数据,在Writer层根据数据关联的文件名动态切换输出文件,避免创建大量Step。
自定义动态Writer实现
public class DynamicFlatFileItemWriter<T> implements ItemWriter<T>, StepExecutionListener { private final Map<String, FlatFileItemWriter<T>> writerCache = new HashMap<>(); private final String baseOutputPath; private final LineAggregator<T> lineAggregator; public DynamicFlatFileItemWriter(String baseOutputPath, LineAggregator<T> lineAggregator) { this.baseOutputPath = baseOutputPath; this.lineAggregator = lineAggregator; } @Override public void write(List<? extends T> items) throws Exception { for (T item : items) { // 假设你的数据对象有getFileName()方法获取目标文件名 String fileName = ((DataItem) item).getFileName(); FlatFileItemWriter<T> writer = writerCache.computeIfAbsent(fileName, this::createWriter); writer.write(Collections.singletonList(item)); } } private FlatFileItemWriter<T> createWriter(String fileName) { FlatFileItemWriter<T> writer = new FlatFileItemWriter<>(); writer.setResource(new FileSystemResource(baseOutputPath + File.separator + fileName)); writer.setLineAggregator(lineAggregator); writer.setAppendAllowed(true); try { writer.afterPropertiesSet(); writer.open(new ExecutionContext()); } catch (Exception e) { throw new RuntimeException("Failed to create writer for file: " + fileName, e); } return writer; } @Override public void beforeStep(StepExecution stepExecution) { // 初始化逻辑 } @Override public ExitStatus afterStep(StepExecution stepExecution) { // 关闭所有缓存的Writer writerCache.values().forEach(writer -> { try { writer.close(); } catch (Exception e) { stepExecution.addFailureException(e); } }); return stepExecution.getExitStatus(); } }
配置Step和Job
@Bean public ItemReader<DataItem> dataItemReader(DataSource dataSource) { JdbcCursorItemReader<DataItem> reader = new JdbcCursorItemReader<>(); reader.setDataSource(dataSource); reader.setSql("SELECT * FROM datas"); reader.setRowMapper((rs, rowNum) -> { DataItem item = new DataItem(); item.setFileName(rs.getString("file_name")); // 其他字段映射 return item; }); return reader; } @Bean public DynamicFlatFileItemWriter<DataItem> dynamicWriter() { LineAggregator<DataItem> lineAggregator = item -> { // 自定义行聚合逻辑,比如CSV格式 return String.join(",", item.getField1(), item.getField2()); }; return new DynamicFlatFileItemWriter<>("/output/path", lineAggregator); } @Bean public Step dataExportStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) { return new StepBuilder("dataExportStep", jobRepository) .<DataItem, DataItem>chunk(100, transactionManager) .reader(dataItemReader(null)) .writer(dynamicWriter()) .listener(dynamicWriter()) .build(); } @Bean public Job dataExportJob(JobRepository jobRepository, Step dataExportStep) { return new JobBuilder("dataExportJob", jobRepository) .start(dataExportStep) .build(); }
方案二:基于Partitioning的并行处理
如果文件数量较多,需要提升处理效率,可以用分区机制,每个分区对应一个文件的处理,支持并行执行。
配置分区Step
// 主Step(分区管理者) @Bean public Step masterStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, Step slaveStep) { return new StepBuilder("masterStep", jobRepository) .partitioner(slaveStep.getName(), filePartitioner()) .step(slaveStep) .gridSize(5) // 并行处理的线程数 .taskExecutor(new SimpleAsyncTaskExecutor()) .build(); } // 分区器:从files表获取文件名作为分区参数 @Bean public Partitioner filePartitioner(DataSource dataSource) { return gridSize -> { Map<String, ExecutionContext> partitions = new HashMap<>(); JdbcTemplate jdbcTemplate = new JdbcTemplate(dataSource); List<String> fileNames = jdbcTemplate.queryForList("SELECT file_name FROM files", String.class); for (String fileName : fileNames) { ExecutionContext context = new ExecutionContext(); context.putString("fileName", fileName); partitions.put("partition-" + fileName, context); } return partitions; }; } // 子Step(每个分区执行的逻辑) @Bean @StepScope public Step slaveStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, @Value("#{stepExecutionContext['fileName']}") String fileName) { return new StepBuilder("slaveStep", jobRepository) .<DataItem, DataItem>chunk(100, transactionManager) .reader(dataItemReaderForFile(fileName)) .writer(fileItemWriter(fileName)) .build(); } // 按文件名过滤的Reader @Bean @StepScope public ItemReader<DataItem> dataItemReaderForFile(DataSource dataSource, @Value("#{stepExecutionContext['fileName']}") String fileName) { JdbcCursorItemReader<DataItem> reader = new JdbcCursorItemReader<>(); reader.setDataSource(dataSource); reader.setSql("SELECT * FROM datas WHERE file_name = ?"); reader.setPreparedStatementSetter(ps -> ps.setString(1, fileName)); reader.setRowMapper((rs, rowNum) -> { DataItem item = new DataItem(); // 字段映射 return item; }); return reader; } // 对应文件的Writer @Bean @StepScope public FlatFileItemWriter<DataItem> fileItemWriter(@Value("#{stepExecutionContext['fileName']}") String fileName) { FlatFileItemWriter<DataItem> writer = new FlatFileItemWriter<>(); writer.setResource(new FileSystemResource("/output/path/" + fileName)); writer.setLineAggregator(item -> String.join(",", item.getField1(), item.getField2())); return writer; }
方案对比
- 动态Writer方案:适合文件数量适中、不需要并行处理的场景,配置简单,单Step即可完成所有文件写入。
- Partitioning方案:适合文件数量多、追求处理效率的场景,支持并行执行,资源利用率更高。
内容的提问来源于stack exchange,提问作者David
相关产品推荐
相关产品推荐

