Spring Batch处理大型CSV文件时,如何从第二个拆分文件开始禁用linesToSkip(1)配置
解决Spring Batch拆分CSV后动态跳过表头的问题
针对你遇到的「仅第一个拆分文件含表头,后续文件需禁用linesToSkip(1)」的问题,我们可以通过收集拆分文件列表+动态判断文件顺序的方式来实现需求,以下是具体的解决方案:
核心思路
- 在拆分文件完成后,收集所有拆分后的文件路径并排序,存入Job Execution Context,确保我们能区分第一个文件和后续文件。
- 使用分区器(Partitioner)遍历所有拆分文件,为每个文件创建独立的处理实例。
- 在
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
相关产品推荐
相关产品推荐

