Spring Batch运行时文件拆分:多服务器实例防重复读取方案
嘿,这个问题刚好是Spring Batch分布式部署里的经典场景——要让多台服务器实例各自处理一部分数据,完全不重复,还要尽量负载均衡。我给你整理了三种实用方案,附带代码示例,你可以根据自己的场景选:
方案一:预拆分源文件(简单直接,适合静态文件)
如果你的源文件是固定不变的,最省心的方式就是提前把大文件拆分成N个(比如5个)等份小文件,每个服务器实例只读取自己对应的小文件。
拆分文件的方式
你可以用Linux的split命令快速拆分:
# 按行数拆分,100万条拆成5份,每份20万条 split -l 200000 large_input.csv split_part_
拆分后会生成split_part_aa、split_part_ab等文件,每个对应20万条记录。
Spring Batch代码配置
让每个实例通过环境变量或启动参数指定要读取的文件:
@Configuration @EnableBatchProcessing public class BatchConfig { // 从启动参数读取当前实例要处理的文件路径 @Value("${batch.input.file}") private String inputFilePath; // 读取文件的ItemReader @Bean public FlatFileItemReader<YourRecord> recordReader() { return new FlatFileItemReaderBuilder<YourRecord>() .name("recordReader") .resource(new FileSystemResource(inputFilePath)) // 配置行映射,适配你的数据格式 .lineMapper(new DefaultLineMapper<YourRecord>() {{ setLineTokenizer(new DelimitedLineTokenizer() {{ setNames("id", "name", "value"); // 替换成你的字段名 }}); setFieldSetMapper(new BeanWrapperFieldSetMapper<YourRecord>() {{ setTargetType(YourRecord.class); }}); }}) .build(); } // 控制台打印的ItemWriter @Bean public ItemWriter<YourRecord> consoleWriter() { return items -> items.forEach(record -> System.out.printf("Instance %s processed: %s%n", System.getenv("INSTANCE_ID"), record.toString()) ); } // 配置Step和Job @Bean public Step processStep() { return stepBuilderFactory.get("processStep") .<YourRecord, YourRecord>chunk(1000) // 批量处理大小,按需调整 .reader(recordReader()) .writer(consoleWriter()) .build(); } @Bean public Job dataProcessJob() { return jobBuilderFactory.get("dataProcessJob") .start(processStep()) .build(); } // 注入Spring Batch的基础工厂类 @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; }
启动各实例
给每个实例指定不同的拆分文件:
# 实例1 java -jar your-batch-app.jar --batch.input.file=./split_part_aa --INSTANCE_ID=1 # 实例2 java -jar your-batch-app.jar --batch.input.file=./split_part_ab --INSTANCE_ID=2 # ...以此类推到实例5
这个方案的优点是零额外依赖、性能高;缺点是文件必须提前拆分,不适合动态生成的文件。
方案二:Spring Batch远程分区(灵活可控,推荐复杂场景)
如果你的场景需要动态分片、故障自动重试,推荐用Spring Batch的远程分区机制。核心逻辑是:
- 一个Master节点负责拆分任务(把100万条记录分成5个分片);
- 多个Worker节点(你的5台服务器)从共享的JobRepository获取分片任务,处理后更新状态;
- 所有节点共享同一个数据库作为JobRepository,保证分片不重复、状态可追踪。
第一步:配置共享JobRepository
所有节点都要连接同一个数据库(比如MySQL),用来存储任务元数据:
@Configuration public class SharedBatchConfig { @Autowired private DataSource dataSource; // 配置共享的JobRepository @Bean public JobRepository jobRepository() throws Exception { JobRepositoryFactoryBean factory = new JobRepositoryFactoryBean(); factory.setDataSource(dataSource); factory.setTransactionManager(new DataSourceTransactionManager(dataSource)); factory.setIsolationLevelForCreate("ISOLATION_REPEATABLE_READ"); factory.setTablePrefix("BATCH_"); // Spring Batch默认表前缀 return factory.getObject(); } // 配置事务管理器 @Bean public PlatformTransactionManager transactionManager() { return new DataSourceTransactionManager(dataSource); } }
第二步:Master节点配置(任务拆分)
Master节点负责把100万条记录拆成5个分片,每个分片对应一个行范围:
@Configuration @EnableBatchProcessing public class MasterConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; @Autowired private JobRepository jobRepository; // 主Job @Bean public Job partitionedJob() throws Exception { return jobBuilderFactory.get("partitionedJob") .start(partitionStep()) .build(); } // 分区Step,负责拆分任务 @Bean public Step partitionStep() throws Exception { return stepBuilderFactory.get("partitionStep") .partitioner("workerStep", filePartitioner()) .partitionHandler(remotePartitionHandler()) .build(); } // 自定义分片器:按行号拆分100万条记录 @Bean public Partitioner filePartitioner() { return gridSize -> { Map<String, ExecutionContext> partitions = new HashMap<>(); long totalRecords = 1000000; long recordsPerPartition = totalRecords / gridSize; for (int i = 0; i < gridSize; i++) { ExecutionContext context = new ExecutionContext(); // 计算当前分片的起止行号(行号从1开始) long startLine = i * recordsPerPartition + 1; long endLine = (i == gridSize - 1) ? totalRecords : (i + 1) * recordsPerPartition; context.putLong("startLine", startLine); context.putLong("endLine", endLine); context.putString("inputFile", "./large_input.csv"); partitions.put("partition-" + i, context); } return partitions; }; } // 远程分区处理器:用消息队列(比如RabbitMQ)分发任务给Worker @Bean public PartitionHandler remotePartitionHandler() { MessageChannelPartitionHandler handler = new MessageChannelPartitionHandler(); handler.setStepName("workerStep"); // Worker节点要执行的Step名称 handler.setJobRepository(jobRepository); handler.setGridSize(5); // 分片数量 // 配置输入输出消息通道(这里需要结合RabbitMQ/Kafka等消息中间件,具体配置参考Spring Batch官方文档) handler.setInputChannel(partitionInputChannel()); handler.setOutputChannel(partitionOutputChannel()); return handler; } // 省略消息通道的配置(需要引入Spring Integration或Spring Cloud Stream依赖) @Bean public MessageChannel partitionInputChannel() { return new DirectChannel(); } @Bean public MessageChannel partitionOutputChannel() { return new DirectChannel(); } }
第三步:Worker节点配置(处理分片任务)
Worker节点根据Master分发的分片参数,只读取指定行范围的记录:
@Configuration @EnableBatchProcessing public class WorkerConfig { @Autowired private StepBuilderFactory stepBuilderFactory; // Worker节点执行的Step @Bean public Step workerStep() { return stepBuilderFactory.get("workerStep") .<YourRecord, YourRecord>chunk(1000) .reader(partitionedRecordReader(null, null, null)) .writer(consoleWriter()) .build(); } // 分片读取器:只处理startLine到endLine之间的记录 @Bean @StepScope // 必须用StepScope才能获取分片上下文参数 public FlatFileItemReader<YourRecord> partitionedRecordReader( @Value("#{stepExecutionContext['startLine']}") Long startLine, @Value("#{stepExecutionContext['endLine']}") Long endLine, @Value("#{stepExecutionContext['inputFile']}") String inputFile) { return new FlatFileItemReaderBuilder<YourRecord>() .name("partitionedRecordReader") .resource(new FileSystemResource(inputFile)) .lineMapper(new DefaultLineMapper<YourRecord>() {{ setLineTokenizer(new DelimitedLineTokenizer() {{ setNames("id", "name", "value"); }}); setFieldSetMapper(new BeanWrapperFieldSetMapper<YourRecord>() {{ setTargetType(YourRecord.class); }}); }}) .recordSeparatorPolicy(new RangeRecordFilter(startLine, endLine)) .build(); } // 自定义行过滤器:只保留指定范围的行 private static class RangeRecordFilter implements RecordSeparatorPolicy { private final long startLine; private final long endLine; private long currentLine = 0; public RangeRecordFilter(long startLine, long endLine) { this.startLine = startLine; this.endLine = endLine; } @Override public boolean isEndOfRecord(String line) { currentLine++; return true; } @Override public String postProcess(String line) { // 只保留起止行范围内的记录 return (currentLine >= startLine && currentLine <= endLine) ? line : null; } @Override public String preProcess(String line) { return line; } } // 控制台打印Writer @Bean public ItemWriter<YourRecord> consoleWriter() { return items -> items.forEach(record -> System.out.printf("Worker %s processed: %s%n", System.getenv("WORKER_ID"), record.toString()) ); } }
注意事项
- 所有节点必须连接同一个JobRepository数据库;
- 需要引入消息中间件依赖(比如RabbitMQ)来实现Master和Worker的通信;
- 如果某个Worker节点失败,Master会重新分发该分片任务,保证数据不丢失。
方案三:轻量分片(无需Master,适合简单场景)
如果不想引入消息中间件,也可以让每个实例根据自身ID自行计算要处理的行范围,不需要中心化的Master节点。
代码配置
@Configuration @EnableBatchProcessing public class LightweightBatchConfig { @Value("${instance.id}") private int instanceId; // 实例ID,从0到4(共5个实例) @Value("${total.instances}") private int totalInstances; // 总实例数,这里是5 @Bean @StepScope public FlatFileItemReader<YourRecord> recordReader() { long totalRecords = 1000000; long recordsPerInstance = totalRecords / totalInstances; // 计算当前实例的起止行号 long startLine = instanceId * recordsPerInstance + 1; long endLine = (instanceId == totalInstances - 1) ? totalRecords : (instanceId + 1) * recordsPerInstance; return new FlatFileItemReaderBuilder<YourRecord>() .name("lightweightRecordReader") .resource(new FileSystemResource("./large_input.csv")) .lineMapper(new DefaultLineMapper<YourRecord>() {{ setLineTokenizer(new DelimitedLineTokenizer() {{ setNames("id", "name", "value"); }}); setFieldSetMapper(new BeanWrapperFieldSetMapper<YourRecord>() {{ setTargetType(YourRecord.class); }}); }}) .recordSeparatorPolicy(new RangeRecordFilter(startLine, endLine)) .build(); } // 省略Writer、Step、Job的配置(和方案一类似) }
启动各实例
# 实例1 java -jar your-batch-app.jar --instance.id=0 --total.instances=5 # 实例2 java -jar your-batch-app.jar --instance.id=1 --total.instances=5 # ...以此类推到实例5
这个方案的优点是轻量无依赖;缺点是没有故障自动重试,某个实例失败需要手动重启。
总结
- 静态文件选方案一,简单高效;
- 复杂场景(动态任务、故障重试)选方案二,功能完整;
- 轻量需求选方案三,快速实现。
内容的提问来源于stack exchange,提问作者reverse
相关产品推荐
相关产品推荐

