You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 15:27:23