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

Spring Batch处理大型CSV文件时,如何从第二个拆分文件开始禁用linesToSkip(1)配置

解决Spring Batch拆分CSV后动态跳过表头的问题

针对你遇到的「仅第一个拆分文件含表头,后续文件需禁用linesToSkip(1)」的问题,我们可以通过收集拆分文件列表+动态判断文件顺序的方式来实现需求,以下是具体的解决方案:


核心思路

  1. 在拆分文件完成后,收集所有拆分后的文件路径并排序,存入Job Execution Context,确保我们能区分第一个文件和后续文件。
  2. 使用分区器(Partitioner)遍历所有拆分文件,为每个文件创建独立的处理实例。
  3. 在FlatFileItemReader中根据当前文件是否为第一个拆分文件,动态设置linesToSkip的数值。

步骤1:修改拆分文件Tasklet,收集拆分后的文件列表

我们需要自定义SystemCommandExecutor,在split命令执行完成后,自动收集所有拆分文件并存储到Job上下文:

@Bean
@StepScope
public SystemCommandTasklet splitFileTasklet(
        @Value("#{jobParameters[filePath]}") final String inputFilePath,
        @Value("#{stepExecution.jobExecution.executionContext}") ExecutionContext executionContext,
        ConfigProperties configProperties) {
    
    SystemCommandTasklet tasklet = new SystemCommandTasklet();
    final File file = BatchUtilities.prefixFile(inputFilePath, AppConstants.PROCESSING_PREFIX);
    // 定义拆分文件的前缀(输入路径+时间戳)
    String splitFilePrefix = configProperties.getBatch().getDataLoad().getInputLocation() + (System.currentTimeMillis() / 1000);
    final String command = configProperties.getBatch().getDataLoadPrep().getSplitCommand() + " " + file.getAbsolutePath() + " " + splitFilePrefix;

    tasklet.setCommand(command);
    tasklet.setTimeout(configProperties.getBatch().getDataLoadPrep().getSplitCommandTimeout());

    // 执行完split命令后,收集所有拆分文件
    tasklet.setSystemCommandExecutor((commandLine, timeout, workingDirectory) -> {
        int exitCode = new DefaultSystemCommandExecutor().executeCommand(commandLine, timeout, workingDirectory);
        if (exitCode == 0) {
            File dir = new File(configProperties.getBatch().getDataLoad().getInputLocation());
            // 筛选出以指定前缀开头的CSV文件
            File[] splitFiles = dir.listFiles((d, name) -> name.startsWith(splitFilePrefix) && name.endsWith(".csv"));
            if (splitFiles != null) {
                // 按文件名排序,保证处理顺序与拆分顺序一致(xaa → xab → xac...)
                Arrays.sort(splitFiles);
                List<String> filePaths = Arrays.stream(splitFiles)
                                              .map(File::getAbsolutePath)
                                              .collect(Collectors.toList());
                // 将文件列表存入Job上下文
                executionContext.put(AppConstants.SPLIT_FILES_LIST, filePaths);
            }
        }
        return exitCode;
    });

    executionContext.put(AppConstants.FILE_PATH_PARAM, file.getPath());
    return tasklet;
}

步骤2:创建文件分区器

使用MultiResourcePartitioner来遍历所有拆分文件,为每个文件分配独立的处理任务:

@Bean
@StepScope
public Partitioner filePartitioner(@Value("#{jobExecutionContext[SPLIT_FILES_LIST]}") List<String> splitFiles) {
    MultiResourcePartitioner partitioner = new MultiResourcePartitioner();
    // 将文件路径转换为Resource数组
    Resource[] resources = splitFiles.stream()
                                     .map(FileSystemResource::new)
                                     .toArray(Resource[]::new);
    partitioner.setResources(resources);
    partitioner.setKeyName("currentFile"); // 用于在Step上下文中标识当前处理的文件
    return partitioner;
}

步骤3:动态配置FlatFileItemReader

修改读取器代码,根据当前文件是否为第一个拆分文件,动态设置linesToSkip:

@Configuration
public class DataLoadReader {
    @Bean
    @StepScope
    public FlatFileItemReader<DemographicData> demographicDataCSVReader(
            @Value("#{stepExecutionContext[currentFile]}") Resource currentFile,
            @Value("#{jobExecutionContext[SPLIT_FILES_LIST]}") List<String> splitFiles) {
        
        // 判断当前文件是否为第一个拆分文件
        boolean isFirstFile = splitFiles.get(0).equals(currentFile.getFile().getAbsolutePath());

        return new FlatFileItemReaderBuilder<DemographicData>()
               .name("data-load-csv-reader")
               .resource(currentFile)
               .linesToSkip(isFirstFile ? 1 : 0) // 第一个文件跳过表头,后续文件不跳过
               .lineMapper(lineMapper())
               .build();
    }

    public LineMapper<DemographicData> lineMapper() {
        DefaultLineMapper<DemographicData> defaultLineMapper = new DefaultLineMapper<>();
        DelimitedLineTokenizer lineTokenizer = new DelimitedLineTokenizer();
        lineTokenizer.setNames("id", "mdl65DecileNum", "mdl66DecileNum", "hhId", "dob", "firstName", "middleName", "lastName", "addressLine1", "addressLine2", "cityName", "stdCode", "zipCode", "zipp4Code", "fipsCntyCd", "fipsStCd", "langName", "regionName", "fipsCntyName", "estimatedIncome");
        defaultLineMapper.setLineTokenizer(lineTokenizer);
        defaultLineMapper.setFieldSetMapper(new DemographicDataFieldSetMapper());
        return defaultLineMapper;
    }
}

步骤4:调整Job配置,使用分区Step

更新Job定义,先执行拆分文件的Step,再执行分区处理的Step:

@Bean
public Job dataLoadJob(JobRepository jobRepository,
                       Step splitFileStep,
                       Step processFileStep) {
    return new JobBuilder("dataLoadJob", jobRepository)
           .start(splitFileStep)
           .next(new StepBuilder("processFilesStep", jobRepository)
                 .partitioner("processFileStep", filePartitioner(null))
                 .step(processFileStep)
                 .gridSize(5) // 根据服务器资源设置并行处理的数量
                 .taskExecutor(taskExecutor())
                 .build())
           .build();
}

@Bean
public TaskExecutor taskExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(3);
    executor.setMaxPoolSize(5);
    executor.setQueueCapacity(10);
    return executor;
}

备选方案(依赖split命令命名规则)

如果不想修改拆分Tasklet,也可以直接通过文件名判断是否为第一个拆分文件(split默认生成的第一个文件后缀为xaa):

// 在Reader中替换判断逻辑
boolean isFirstFile = currentFile.getFilename().endsWith("xaa.csv");

⚠️ 注意:这种方式依赖split命令的默认命名规则,若后续修改split的参数(如--suffix-length)则会失效,建议优先使用第一种方案。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 22:42:32