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

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的远程分区机制。核心逻辑是:

  1. 一个Master节点负责拆分任务(把100万条记录分成5个分片);
  2. 多个Worker节点(你的5台服务器)从共享的JobRepository获取分片任务,处理后更新状态;
  3. 所有节点共享同一个数据库作为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 10:36:09