Spring Batch批处理案例

wen java案例 2

本文目录导读:

Spring Batch批处理案例

  1. 目录导读
  2. 为什么需要Spring Batch?——批处理的痛点与选型考量
  3. 核心架构解剖:Job、Step、Chunk模型如何运作?
  4. 实战案例:银行日终对账批处理系统
  5. 性能调优三把斧:分区、并行、异步
  6. 常见踩坑与面试问答(含可靠性保障)
  7. 总结与演进建议

Spring Batch批处理案例实战:从零搭建高效银行对账系统(附完整代码)

目录导读

  1. 为什么需要Spring Batch?——批处理的痛点与选型考量
  2. 核心架构解剖:Job、Step、Chunk模型如何运作?
  3. 实战案例:银行日终对账批处理系统(需求分析→代码实现)
  4. 性能调优三把斧:分区、并行、异步
  5. 常见踩坑与面试问答(含可靠性保障)

为什么需要Spring Batch?——批处理的痛点与选型考量

在企业级应用中,我们经常遇到定时、大批量、无交互的数据处理场景,每日报表生成、数据迁移、对账清算,如果直接用for循环 + JDBC硬编码,你会面临三大问题:

  • 无状态管理:程序中途崩溃,无法从断点续跑
  • 无重试机制:一行脏数据导致整批任务失败
  • 性能瓶颈:单线程处理千万级数据,耗时以小时计

Spring Batch 作为Spring全家桶的批处理框架,完美解决了上述问题,它提供了声明式作业编排事务性分块处理跳过/重试机制监控与重启等核心能力,且与Spring Boot深度集成,是Java生态中最主流的批处理解决方案。

选型对比:相比自研线程池 + 循环,Spring Batch学习曲线略陡峭,但换来的是生产级的健壮性,如果你的场景只是简单ETL,可以考虑Spring Batch + Spring Integration简化版;若涉及复杂分布式调度,可结合xxl-jobQuartz管理触发。


核心架构解剖:Job、Step、Chunk模型如何运作?

在进入案例前,我们先澄清三个关键概念:

  • Job(作业):一个完整的批处理过程,由1个或多个Step组成,对账作业”包含“读取银行流水”和“生成差异报告”两个Step。
  • Step(步骤):作业中的一个独立阶段,每个Step默认采用Chunk-oriented处理(面向块)。
  • Chunk(块):任务划分为若干“块”,每个块内执行:read(读若干条)→ process(处理)→ write(批量写),然后提交一次事务。块大小(commit-interval) 直接影响性能与事务粒度。

执行时序图(简化):

JobLauncher → Job → Step(1..n) → Chunk(读→处理→写) → Repeat

实战案例:银行日终对账批处理系统

1 需求分析

假设我们有一个银行核心系统,每天凌晨2点需要完成:

  • 读取:从FTP下载当日第三方支付机构交易流水(CSV格式,约100万条)
  • 处理:与本地数据库的银行交易记录比对,标记差异状态(匹配/金额不一致/缺失)
  • 输出:生成差异报表(TXT),并更新数据库对账状态字段

2 环境与技术栈

  • Spring Boot 2.7.x
  • Spring Batch 4.3.x
  • MyBatis-Plus(数据访问)
  • H2内存数据库(测试用,生产可换MySQL)
  • Lombok,Maven

3 核心代码实践

步骤1:添加依赖

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-batch</artifactId>
</dependency>

步骤2:定义实体与Mapper

// 第三方流水表 DTO
@Data
public class ThirdPartyFlow {
    private String txId;      // 交易号
    private BigDecimal amount; // 金额
    private Date txDate;      // 交易日期
}

步骤3:配置Job与Step(核心)

@Configuration
public class BatchConfig {
    @Autowired
    private JobBuilderFactory jobBuilderFactory;
    @Autowired
    private StepBuilderFactory stepBuilderFactory;
    @Autowired
    private DataSource dataSource;
    // 1. 定义Job:对账作业
    @Bean
    public Job reconcileJob() {
        return jobBuilderFactory.get("reconcileJob")
                .start(reconcileStep())
                .build();
    }
    // 2. 定义Step:采用Chunk模型,每500条提交一次事务
    @Bean
    public Step reconcileStep() {
        return stepBuilderFactory.get("reconcileStep")
                .<ThirdPartyFlow, ReconcileResult>chunk(500)
                .reader(flowItemReader())          // 读取CSV
                .processor(compareProcessor())     // 比对逻辑
                .writer(resultItemWriter())        // 写差异报告并更新状态
                .faultTolerant()
                .skipLimit(100)                    // 允许跳过100条异常数据
                .skip(ParseException.class)        // 解析异常跳过
                .retryLimit(3)                     // 重试3次
                .retry(DataAccessException.class)  // 数据库异常重试
                .build();
    }
    // 3. 读取器:从CSV读取(FlatFileItemReader)
    @Bean
    public FlatFileItemReader<ThirdPartyFlow> flowItemReader() {
        return new FlatFileItemReaderBuilder<ThirdPartyFlow>()
                .name("flowItemReader")
                .resource(new FileSystemResource("/data/third_party_flow.csv"))
                .delimited()
                .names("txId", "amount", "txDate")
                .fieldSetMapper(new BeanWrapperFieldSetMapper<>() {{
                    setTargetType(ThirdPartyFlow.class);
                }})
                .build();
    }
    // 4. 处理器:与本地交易比对
    @Bean
    public ItemProcessor<ThirdPartyFlow, ReconcileResult> compareProcessor() {
        return thirdPartyFlow -> {
            // 伪代码:调用Mapper查询本地交易
            LocalTrade localTrade = localTradeMapper.selectByTxId(thirdPartyFlow.getTxId());
            ReconcileResult result = new ReconcileResult();
            if (localTrade == null) {
                result.setOddType("缺失本地记录");
            } else if (localTrade.getAmount().compareTo(thirdPartyFlow.getAmount()) != 0) {
                result.setOddType("金额不一致");
            } else {
                result.setOddType("匹配正常");
            }
            result.setTxId(thirdPartyFlow.getTxId());
            return result;
        };
    }
    // 5. 写入器:输出差异文件 & 更新数据库状态
    @Bean
    public ItemWriter<ReconcileResult> resultItemWriter() {
        return list -> {
            // 写入TXT文件
            try (BufferedWriter writer = Files.newBufferedWriter(Paths.get("/data/result_"
                    + System.currentTimeMillis() + ".txt"))) {
                for (ReconcileResult r : list) {
                    if (!"匹配正常".equals(r.getOddType())) {
                        writer.write(r.getTxId() + "|" + r.getOddType() + "\n");
                    }
                }
            }
            // 批量更新数据库状态(使用MyBatis-Plus)
            reconcileResultMapper.batchUpdateStatus(list);
        };
    }
}

步骤4:启动作业(通过CommandLineRunner)

@Component
public class JobLaunchRunner implements CommandLineRunner {
    @Autowired
    private JobLauncher jobLauncher;
    @Autowired
    private Job reconcileJob;
    @Override
    public void run(String... args) throws Exception {
        JobParameters params = new JobParametersBuilder()
                .addDate("date", new Date())
                .toJobParameters();
        jobLauncher.run(reconcileJob, params);
    }
}

4 运行与验证

  • 启动Spring Boot应用,观察控制台日志:Job: [SimpleJob: [reconcileJob]] completed successfully
  • 检查/data/result_xxx.txt是否生成了差异记录
  • 数据库对账状态字段已更新

性能调优三把斧:分区、并行、异步

当数据量从百万级涨到亿级,单机默认Chunk处理效率不够,以下三种优化手段可叠加使用:

  1. 分区(Partitioning):将数据按“日期”或“银行编号”拆分为多个区(Partition),每个区由独立的Step执行,示例:

    stepBuilderFactory.get("partitionStep")
     .partitioner("workerStep", new CustomPartitioner())
     .taskExecutor(new SimpleAsyncTaskExecutor()) // 并行执行
     .build();
  2. 多线程Step:设置taskExecutor,让一个Step内部使用线程池并行处理多个Chunk,注意保证Reader的线程安全,推荐使用SynchronizedItemStreamReader

  3. 远程分块与异步处理:通过MessageChannel将数据发送给Kafka或RabbitMQ,由远程消费者完成处理,适合分布式架构。


常见踩坑与面试问答(含可靠性保障)

Q1:Job中途宕机了,如何续跑? A:Spring Batch通过JobRepository(默认存在数据库)持久化执行上下文(JobExecutionContext),重启后,使用相同的JobParameters(如带日期参数)启动,框架会检测到已存在的JobInstance,默认不重复执行,若需强制重跑,需新增JobParameters(如加一个时间戳参数);若想从失败的Step处续跑,配置Jobstart(...).on("FAILED").to(...).end()恢复流程,或使用JobOperatorrestart()

Q2:为什么我的ItemWriter没有事务? A:Spring Batch的Chunk事务是由框架管理的,默认每个Chunk提交一次事务,但请注意:Writer必须在事务性资源(如JdbcTemplate)上操作才有事务保护,如果你直接操作文件输出,文件操作不是事务性的,建议将“写文件”放在StepExecutionListener.afterChunk回调中,而Writer只负责数据库更新。

Q3:数据量大,内存溢出怎么办? A:调大commit-interval不是万能药,建议:

  • 使用游标式Reader(JdbcCursorItemReader)而非分页式,避免一次性加载大量实体到内存。
  • 若必须用分页,必须配置saveState=true并设置合理的pageSize(如1000)。
  • 考虑分区+并行,分散内存压力。

Q4:如何处理文件乱码/脏数据? AFlatFileItemReader支持LineMapper自定义解析;在Processor中捕获ParseException,配合skip机制,更好的做法是写一条SkipListener,把异常数据单独记录到日志表,便于人工处理。

Q5:Spring Batch与Quartz如何配合? A:Spring Batch只负责“执行批处理”,不负责“定时触发”,常用方案:Quartz调度器定时调用JobLauncher.run(),Spring Boot 2.x还支持@Scheduled注解,配合JobLauncher即可实现简单的定时批处理。


总结与演进建议

通过上述案例,我们完整实现了“文件读取→业务比对→结果输出”的批处理闭环,Spring Batch的核心价值在于可靠性与可维护性,而非单纯的性能——它帮你管理了状态、重试、跳过、监控等繁琐细节。

进阶方向

  • 若需要可视化监控,集成Spring Batch AdminMicrometer + Prometheus/Grafana
  • 若批处理任务非常复杂(如多数据源、动态流程),可研究Batch ScalingSpring Cloud Task
  • 注意:批处理不等于流处理,若实时性要求高(毫秒级),请转向Kafka StreamsFlink

建议实践清单

  1. 先跑通上述Demo,理解Job/Step/Chunk的生命周期。
  2. 将你的现有循环代码改造为Batch框架,体验异常恢复的便利。
  3. 压测调整commit-interval,找到性能拐点。

希望这篇文章能帮你告别“脚本式批处理”,走向工程化、可治理的批处理之路,如有具体问题,欢迎在评论区交流。


(注:本文所有代码均基于Spring Boot 2.7.x与Spring Batch 4.3.x编写,生产环境请根据具体版本调整API。)

抱歉,评论功能暂时关闭!