乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 批处理概述与 Spring Batch 核心理念

    • Spring Batch 概述
    • Job-Instance-Execution 三层生命周期
    • 三层架构
    • Chunk 处理模型
    • Tasklet 处理模型
  • 第2章 Job与作业配置

    • Job 详解
    • Job 配置
    • JobLauncher 详解
    • JobParameters 详解
    • JobRepository 详解
  • 第3章 Step与执行模型

    • Step 详解
    • Step 配置
    • ExecutionContext 详解
    • ItemStream 与状态管理
    • StepScope 与 JobScope
  • 第4章 ItemReader数据读取

    • ItemReader 详解
    • FlatFileItemReader 详解
    • JdbcItemReader 详解
    • MultiResourceItemReader 详解
  • 第5章 ItemProcessor 数据处理

    • ItemProcessor 详解
  • 第6章 ItemWriter 数据写出

    • ItemWriter 详解
    • FlatFileItemWriter 详解
    • CompositeItemWriter 详解
    • JdbcBatchItemWriter 详解
  • 第7章 Chunk 处理与事务边界

    • Chunk 处理模型详解
    • 事务边界
  • 第8章 作业参数与启动

    • 命令行与 Web 启动
  • 第9章 监听器与拦截器

    • 监听器详解
  • 第10章 重启重试与跳过策略

    • 重启与重试
    • 跳过策略
  • 第11章 作业流与条件决策

    • 作业流与条件决策
  • 第12章 分区与并行处理

    • 分区详解
    • 并行 Step 与 Split
  • 第13章 与 Spring Boot 集成实践

    • Spring Boot 集成
  • 第15章 运维与监控

    • 大作业设计模式
    • 运维与监控

大作业设计模式

本章定位:以生产环境视角,掌握驱动表 + 分批提交、增量处理和断点续传三种核心设计模式,将 Spring Batch 知识体系转化为工程实践。

定义与作用

Spring Batch 提供了基础 API,但大型生产级批处理需要更高层的设计模式来保证性能、可靠性和可维护性。三种核心模式覆盖了 90% 的企业批处理需求。

飞翔科技架构师白歌的总结:

"批处理的本质就一件事:用可控的成本在可控的时间内安全地把 A 点的数据搬到 B 点。三种设计模式分别解决三个核心问题——全量怎么搬、增量怎么搬、出错了怎么续。"

模式一:驱动表 + 分批提交

适用场景:百万级以上数据的全量 ETL 处理。

@Bean
public Step fullEtlStep(JobRepository jobRepository,
                          PlatformTransactionManager tx,
                          DataSource dataSource) {
    JdbcPagingItemReader<SourceRecord> reader = new JdbcPagingItemReaderBuilder<SourceRecord>()
            .name("fullReader")
            .dataSource(dataSource)
            .selectClause("SELECT *")
            .fromClause("FROM source_table")
            .sortKeys(Map.of("id", Order.ASCENDING))
            .pageSize(1000)           // 大分页,减少 SQL 次数
            .rowMapper(new BeanPropertyRowMapper<>(SourceRecord.class))
            .saveState(true)
            .build();

    return new StepBuilder("fullEtlStep", jobRepository)
            .<SourceRecord, TargetRecord>chunk(500, tx) // 大 Chunk
            .reader(reader)
            .processor(item -> transform(item))
            .writer(new JdbcBatchItemWriterBuilder<TargetRecord>()
                    .dataSource(dataSource)
                    .sql("INSERT INTO target_table VALUES (?, ?, ?)")
                    .itemPreparedStatementSetter((item, ps) -> {
                        ps.setLong(1, item.getId());
                        ps.setString(2, item.getName());
                        ps.setBigDecimal(3, item.getValue());
                    })
                    .build())
            .build();
}

性能指标(500 万条,MySQL 8.0):

Chunk SizePage Size耗时事务数
100100~180s50000
500500~60s10000
10001000~40s5000

模式二:增量处理

适用场景:定时同步增量数据,只处理自上次运行以来的变更。

@Bean
@StepScope
public JdbcPagingItemReader<Order> incrementalReader(
        DataSource dataSource,
        @Value("#{jobParameters['last.run.time']}") String lastRunTime) {

    Map<String, Object> params = new HashMap<>();
    params.put("lastRunTime", LocalDateTime.parse(lastRunTime));

    return new JdbcPagingItemReaderBuilder<Order>()
            .name("incrementalReader")
            .dataSource(dataSource)
            .selectClause("SELECT *")
            .fromClause("FROM orders")
            .whereClause("WHERE updated_at > :lastRunTime")
            .parameterValues(params)
            .sortKeys(Map.of("updated_at", Order.ASCENDING))
            .pageSize(500)
            .rowMapper(new BeanPropertyRowMapper<>(Order.class))
            .saveState(true)
            .build();
}

⚠️ 增量处理的 Job 参数:每次运行使用不同的 last.run.time 参数,创建一个新的 JobInstance。

模式三:断点续传 + 幂等写入

适用场景:需保证数据不重不丢的关键任务。

@Bean
public Step idempotentStep(JobRepository jobRepository,
                              PlatformTransactionManager tx,
                              DataSource dataSource) {
    return new StepBuilder("idempotentStep", jobRepository)
            .<Order, Order>chunk(500, tx)
            .reader(new JdbcPagingItemReaderBuilder<Order>()
                    .name("idempotentReader")
                    .dataSource(dataSource)
                    .selectClause("SELECT *")
                    .fromClause("FROM orders")
                    .sortKeys(Map.of("id", Order.ASCENDING))
                    .pageSize(500)
                    .rowMapper(new BeanPropertyRowMapper<>(Order.class))
                    .saveState(true)          // 保存分页位置
                    .build())
            .writer(new JdbcBatchItemWriterBuilder<Order>()
                    .dataSource(dataSource)
                    .sql('INSERT INTO orders_processed (order_id, amount, processed_at) '
                        'VALUES (:orderId, :amount, NOW()) '
                        'ON DUPLICATE KEY UPDATE amount = VALUES(amount)')
                    .itemSqlParameterSourceProvider(
                        new BeanPropertyItemSqlParameterSourceProvider<>())
                    .build())
            .build();
}

三种模式对比:

模式核心机制适用数据量关键风险
驱动表 + 分批分页 Reader + 大 Chunk全量 100 万+事务时间长,锁竞争
增量处理WHERE updated_at > :lastRun增量 1 万 - 100 万需可靠的时间戳字段
断点续传 + 幂等saveState + ON DUPLICATE KEY任意规模需要排序键和唯一键

易错场景与避坑

反例一:增量模式但 Job 参数不变

// ❌ 每次都传相同的 run.id → 同一个 JobInstance → 第二次报 AlreadyComplete
JobParameters params = new JobParametersBuilder()
        .addString("run.id", "fixed-value") // ❌ 固定值!
        .toJobParameters();

正确做法:每次运行用时间戳或递增 ID 作为 identifying 参数。

反例二:全量处理不使用 saveState

// ❌ saveState=false → 500 万条处理到 450 万条崩溃 → 从头重来
reader.setSaveState(false);

面试高频考点

Q1:如何设计一个既要全量又要增量的批处理系统?

全量 Job 每周日凌晨运行一次(用驱动表模式),增量 Job 每 15 分钟运行一次(用 updated_at > :lastRunTime 模式)。全量和增量使用不同的 Job 定义,互不干扰。

Q2:三种模式能否组合使用?

可以。例如全量 ETL 用「驱动表 + 分批提交」模式保证性能,同时 Reader 实现 ItemStream 支持断点续传,Writer 用 ON DUPLICATE KEY UPDATE 保证幂等。三种模式不是互斥的,而是正交增强的。


上一章:运维与监控全文完 | Spring Batch 独立教程

下一页
运维与监控