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

Spring Batch:基于外部API响应前置过滤CSV数据方案咨询

Spring Batch作业启动前调用外部API过滤CSV数据实现方案

针对你的需求,这里提供一套完整的实现方案,核心是通过自定义前置Tasklet完成API调用和CSV过滤,调整作业执行顺序后实现数据预过滤,具体如下:

整体流程

作业执行顺序调整为:

  1. 调用外部API获取允许处理的数据集(白名单)
  2. 过滤原始CSV文件,生成仅包含白名单数据的临时CSV
  3. 使用FlatFileItemReader读取临时CSV,执行后续业务Tasklet/处理步骤

具体实现步骤

1. 自定义前置Tasklet:API调用+CSV过滤

这个Tasklet负责完成API调用获取白名单,以及过滤原始CSV生成临时文件的核心逻辑:

@Component
public class PreProcessingTasklet implements Tasklet {

    private final RestTemplate restTemplate;
    private final ResourceLoader resourceLoader;

    // 构造注入依赖
    public PreProcessingTasklet(RestTemplate restTemplate, ResourceLoader resourceLoader) {
        this.restTemplate = restTemplate;
        this.resourceLoader = resourceLoader;
    }

    @Override
    public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception {
        // 1. 调用外部API获取允许处理的ID集合(根据实际API响应结构调整)
        Set<String> allowedIds = fetchAllowedIdsFromApi();

        // 2. 读取原始CSV并过滤生成临时文件
        Resource originalCsv = resourceLoader.getResource("classpath:input/source_data.csv");
        File tempFilteredCsv = File.createTempFile("batch_filtered_", ".csv");
        
        // 用缓冲流处理文件,避免内存溢出
        try (BufferedReader reader = new BufferedReader(new FileReader(originalCsv.getFile()));
             BufferedWriter writer = new BufferedWriter(new FileWriter(tempFilteredCsv))) {
            // 保留CSV表头
            String headerLine = reader.readLine();
            writer.write(headerLine);
            writer.newLine();

            // 逐行过滤数据(假设CSV第一列为ID,根据实际结构调整)
            String dataLine;
            while ((dataLine = reader.readLine()) != null) {
                String[] fields = dataLine.split(",");
                if (allowedIds.contains(fields[0].trim())) { // 注意去除字段前后空格
                    writer.write(dataLine);
                    writer.newLine();
                }
            }
        }

        // 将临时文件路径存入作业上下文,供后续Reader读取
        chunkContext.getStepContext().getJobExecutionContext()
            .put("filteredCsvPath", tempFilteredCsv.getAbsolutePath());
        return RepeatStatus.FINISHED;
    }

    private Set<String> fetchAllowedIdsFromApi() {
        // 替换为实际API地址和响应解析逻辑
        ResponseEntity<List<String>> apiResponse = restTemplate.exchange(
                "https://your-external-api.com/allowed-data",
                HttpMethod.GET,
                null,
                new ParameterizedTypeReference<List<String>>() {}
        );
        return new HashSet<>(apiResponse.getBody());
    }
}

2. 配置Batch作业流程

调整作业步骤顺序,将前置处理Tasklet作为第一个执行步骤,后续步骤读取临时CSV文件:

@Configuration
@EnableBatchProcessing
public class BatchJobConfig {

    private final JobBuilderFactory jobBuilderFactory;
    private final StepBuilderFactory stepBuilderFactory;
    private final PreProcessingTasklet preProcessingTasklet;

    public BatchJobConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, PreProcessingTasklet preProcessingTasklet) {
        this.jobBuilderFactory = jobBuilderFactory;
        this.stepBuilderFactory = stepBuilderFactory;
        this.preProcessingTasklet = preProcessingTasklet;
    }

    // 前置处理步骤:API调用+CSV过滤
    @Bean
    public Step preProcessingStep() {
        return stepBuilderFactory.get("preProcessingStep")
                .tasklet(preProcessingTasklet)
                .build();
    }

    // 核心数据处理步骤(复用原有Reader/Processor/Writer,仅调整Reader的文件来源)
    @Bean
    public Step dataProcessingStep() {
        return stepBuilderFactory.get("dataProcessingStep")
                .<YourDataModel, YourDataModel>chunk(100) // 调整chunk大小适配你的数据量
                .reader(filteredCsvReader(null))
                .processor(yourItemProcessor())
                .writer(yourItemWriter())
                .build();
    }

    // StepScope注解实现动态读取作业上下文的临时文件路径
    @Bean
    @StepScope
    public FlatFileItemReader<YourDataModel> filteredCsvReader(
            @Value("#{jobExecutionContext['filteredCsvPath']}") String filteredFilePath) {
        return new FlatFileItemReaderBuilder<YourDataModel>()
                .name("filteredCsvReader")
                .resource(new FileSystemResource(filteredFilePath))
                .delimited()
                .names("id", "name", "value") // 替换为你的CSV实际字段名
                .targetType(YourDataModel.class)
                .build();
    }

    // 原有业务处理逻辑:ItemProcessor
    @Bean
    public ItemProcessor<YourDataModel, YourDataModel> yourItemProcessor() {
        return item -> {
            // 你的业务处理逻辑,比如数据转换、校验
            return item;
        };
    }

    // 原有业务输出逻辑:ItemWriter
    @Bean
    public ItemWriter<YourDataModel> yourItemWriter() {
        return items -> {
            // 你的数据输出逻辑,比如写入数据库、调用其他服务
        };
    }

    // 组装作业流程:先执行前置处理,再执行核心数据处理
    @Bean
    public Job dataProcessingJob() {
        return jobBuilderFactory.get("dataProcessingJob")
                .start(preProcessingStep())
                .next(dataProcessingStep())
                .build();
    }
}

关键细节优化

1. 临时文件清理

作业执行完成后,建议添加清理步骤删除临时CSV文件,避免磁盘占用:

@Component
public class TempFileCleanupTasklet implements Tasklet {

    @Override
    public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception {
        String tempFilePath = (String) chunkContext.getStepContext().getJobExecutionContext().get("filteredCsvPath");
        File tempFile = new File(tempFilePath);
        if (tempFile.exists()) {
            tempFile.delete();
        }
        return RepeatStatus.FINISHED;
    }
}

然后在作业配置中添加该步骤:

@Bean
public Step cleanupStep() {
    return stepBuilderFactory.get("cleanupStep")
            .tasklet(tempFileCleanupTasklet)
            .build();
}

// 调整作业流程
@Bean
public Job dataProcessingJob() {
    return jobBuilderFactory.get("dataProcessingJob")
            .start(preProcessingStep())
            .next(dataProcessingStep())
            .next(cleanupStep())
            .build();
}

2. API异常处理

在fetchAllowedIdsFromApi方法中添加异常捕获,避免API调用失败导致作业直接中断:

private Set<String> fetchAllowedIdsFromApi() throws RuntimeException {
    try {
        ResponseEntity<List<String>> apiResponse = restTemplate.exchange(...);
        return new HashSet<>(apiResponse.getBody());
    } catch (RestClientException e) {
        // 根据业务需求选择:抛出异常终止作业,或返回空集合(不处理任何数据),或返回全量允许集合
        throw new RuntimeException("调用外部API失败,作业终止", e);
    }
}

3. 复杂CSV解析

如果CSV包含带逗号的字段(如"Doe, John"),不要直接用split(","),建议使用OpenCSV等专业库解析:

// 引入OpenCSV依赖
<dependency>
    <groupId>com.opencsv</groupId>
    <artifactId>opencsv</artifactId>
    <version>5.6</version>
</dependency>

然后在Tasklet中替换解析逻辑:

try (CSVReader reader = new CSVReader(new FileReader(originalCsv.getFile()));
     CSVWriter writer = new CSVWriter(new FileWriter(tempFilteredCsv))) {
    String[] header = reader.readNext();
    writer.writeNext(header);
    String[] line;
    while ((line = reader.readNext()) != null) {
        if (allowedIds.contains(line[0].trim())) {
            writer.writeNext(line);
        }
    }
}

可选优化方案

  • 内存过滤替代文件IO:如果CSV文件较小,可直接将过滤后的数据存入作业上下文,用ListItemReader读取内存数据,避免磁盘IO开销。
  • API响应缓存:如果API响应不会频繁更新,可添加Redis或本地缓存,避免每次作业启动都调用API,提升执行效率。

内容的提问来源于stack exchange,提问作者Chaithra Rai

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 00:35:24