如何用Spring Batch实现多文件读取处理后原地更新写入?
实现方案说明
Spring Batch完全支持读取多文件、修改行数据后写回原文件的场景,核心思路是让每个处理后的Item关联其来源文件路径,通过ClassifierCompositeItemWriter按文件路径分类写入。以下是具体实现步骤和代码示例:
1. 定义带源文件标识的Item类
需要给每行数据附加来源文件的路径,用于后续分类写入:
public class ProcessedLine { // 原始/修改后的行内容 private String content; // 该行数据对应的源文件绝对路径 private String sourceFilePath; // Getter & Setter public String getContent() { return content; } public void setContent(String content) { this.content = content; } public String getSourceFilePath() { return sourceFilePath; } public void setSourceFilePath(String sourceFilePath) { this.sourceFilePath = sourceFilePath; } }
2. 实现资源感知的ItemReader
包装MultiResourceItemReader,在读取每行数据时自动附加当前源文件路径:
public class ResourceAwareItemReader implements ItemReader<ProcessedLine> { private final MultiResourceItemReader<String> delegate; private String currentFilePath; public ResourceAwareItemReader(MultiResourceItemReader<String> delegate) { this.delegate = delegate; } @Override public ProcessedLine read() throws Exception { String line = delegate.read(); if (line == null) { currentFilePath = null; return null; } // 首次读取当前文件时,记录文件路径 if (currentFilePath == null) { Resource currentResource = delegate.getCurrentResource(); currentFilePath = currentResource.getFile().getAbsolutePath(); } ProcessedLine processedLine = new ProcessedLine(); processedLine.setContent(line); processedLine.setSourceFilePath(currentFilePath); return processedLine; } }
3. 自定义行数据处理器
按业务逻辑修改行内容,这里以字符串替换为例:
public class LineProcessingItemProcessor implements ItemProcessor<ProcessedLine, ProcessedLine> { @Override public ProcessedLine process(ProcessedLine item) throws Exception { // 自定义修改逻辑,比如替换指定内容 String modifiedContent = item.getContent().replace("old_value", "new_value"); item.setContent(modifiedContent); return item; } }
4. 实现分类写入的ItemWriter
使用ClassifierCompositeItemWriter,根据Item的源文件路径匹配对应的FlatFileItemWriter。为避免原文件损坏,先写入临时文件,Job完成后再替换原文件:
4.1 实现文件路径分类器
public class FilePathClassifier implements Classifier<ProcessedLine, ItemWriter<? super ProcessedLine>> { private final Map<String, ItemWriter<ProcessedLine>> writerMap = new HashMap<>(); private final Map<String, String> tempFileMap = new HashMap<>(); @Override public ItemWriter<? super ProcessedLine> classify(ProcessedLine item) { String filePath = item.getSourceFilePath(); return writerMap.computeIfAbsent(filePath, this::createTempFileWriter); } // 创建临时文件Writer private FlatFileItemWriter<ProcessedLine> createTempFileWriter(String originalFilePath) { try { File originalFile = new File(originalFilePath); // 在原文件同目录创建临时文件 File tempFile = File.createTempFile("batch_temp_", ".tmp", originalFile.getParentFile()); tempFileMap.put(originalFilePath, tempFile.getAbsolutePath()); FlatFileItemWriter<ProcessedLine> writer = new FlatFileItemWriter<>(); writer.setResource(new FileSystemResource(tempFile)); // 将ProcessedLine的content字段作为行内容写入 writer.setLineAggregator(ProcessedLine::getContent); writer.setAppend(false); // 覆盖临时文件(每次写入新内容) writer.afterPropertiesSet(); return writer; } catch (Exception e) { throw new RuntimeException("初始化临时文件Writer失败: " + originalFilePath, e); } } // Job完成后替换原文件 public void replaceOriginalFiles() throws IOException { for (Map.Entry<String, String> entry : tempFileMap.entrySet()) { File original = new File(entry.getKey()); File temp = new File(entry.getValue()); // 替换原文件,若原文件存在则覆盖 Files.move(temp.toPath(), original.toPath(), StandardCopyOption.REPLACE_EXISTING); } } // Job失败时清理临时文件 public void cleanTempFiles() { tempFileMap.values().forEach(tempPath -> new File(tempPath).delete()); } }
4.2 配置Job执行监听器
确保Job成功完成后替换原文件,失败则清理临时文件:
public class FileReplacementListener implements JobExecutionListener { private final FilePathClassifier classifier; public FileReplacementListener(FilePathClassifier classifier) { this.classifier = classifier; } @Override public void beforeJob(JobExecution jobExecution) {} @Override public void afterJob(JobExecution jobExecution) { if (jobExecution.getStatus() == BatchStatus.COMPLETED) { try { classifier.replaceOriginalFiles(); } catch (IOException e) { jobExecution.setStatus(BatchStatus.FAILED); throw new RuntimeException("替换原文件失败", e); } } else { classifier.cleanTempFiles(); } } }
5. 完整Batch配置类
@Configuration @EnableBatchProcessing public class BatchConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; // 配置多文件读取器 @Bean public MultiResourceItemReader<String> multiResourceItemReader() { MultiResourceItemReader<String> reader = new MultiResourceItemReader<>(); // 指定读取的文件路径,这里读取input目录下所有txt文件 reader.setResources(new Resource[]{new FileSystemResource("input/*.txt")}); // 内部行读取器,直接读取每行字符串 FlatFileItemReader<String> lineReader = new FlatFileItemReader<>(); lineReader.setLineMapper((line, lineNum) -> line); reader.setDelegate(lineReader); return reader; } @Bean public ItemReader<ProcessedLine> resourceAwareItemReader() { return new ResourceAwareItemReader(multiResourceItemReader()); } @Bean public ItemProcessor<ProcessedLine, ProcessedLine> lineProcessor() { return new LineProcessingItemProcessor(); } @Bean public FilePathClassifier filePathClassifier() { return new FilePathClassifier(); } @Bean public ClassifierCompositeItemWriter<ProcessedLine> compositeWriter() { ClassifierCompositeItemWriter<ProcessedLine> writer = new ClassifierCompositeItemWriter<>(); writer.setClassifier(filePathClassifier()); return writer; } @Bean public Step processingStep() { return stepBuilderFactory.get("processingStep") .<ProcessedLine, ProcessedLine>chunk(100) // 按100行一批处理 .reader(resourceAwareItemReader()) .processor(lineProcessor()) .writer(compositeWriter()) .build(); } @Bean public Job fileProcessingJob() { return jobBuilderFactory.get("fileProcessingJob") .start(processingStep()) .listener(new FileReplacementListener(filePathClassifier())) .build(); } }
关键注意事项
- 文件读写安全:通过临时文件中转,避免直接写入原文件导致数据损坏。
- 内存占用:根据文件大小调整chunk参数,避免OOM。
- 异常处理:监听器会在Job失败时自动清理临时文件,可根据业务需求扩展异常处理逻辑。
内容的提问来源于stack exchange,提问作者breakingcode
相关产品推荐
相关产品推荐

