Article

Spring Batch

更新于:2026-07-14

第1章 Spring Batch 概述与核心概念

1.1 什么是 Spring Batch

概念名称说明注意事项
Spring Batch一个轻量级、全面的批处理框架,用于开发强大的批处理应用程序。适用于处理大量数据的场景,不适合实时或低延迟处理。
批处理(Batch Processing)指在无用户交互的情况下,自动执行一系列任务来处理大量数据。通常在后台运行,强调吞吐量和资源利用率。
开源项目由 Spring 社区维护,基于 Spring Framework 构建,与 Spring 生态无缝集成。需要 Java 环境支持,推荐使用 Java 8 或更高版本。

1.2 Spring Batch 的应用场景

应用场景说明注意事项
数据迁移将旧系统数据迁移到新系统,如数据库迁移、文件格式转换等。需确保数据一致性,支持断点续传和错误处理。
报表生成定期生成财务报表、统计报告等批量输出任务。通常在夜间或低峰期执行,避免影响在线业务。
批量交易处理处理银行对账、订单结算、工资发放等周期性任务。要求高可靠性,支持事务管理和回滚机制。
ETL(Extract-Transform-Load)从多个源提取数据,进行清洗转换后加载到目标系统(如数据仓库)。强调数据质量控制和处理效率。
文件导入导出批量导入 CSV、XML 文件到数据库,或从数据库导出为文件。需处理大文件分页读取,防止内存溢出。

1.3 Spring Batch 核心架构组件

组件名称说明注意事项
Job批处理任务的最外层容器,代表一个完整的业务流程。一个 Job 可包含多个 Step,Job 是可重启的单位。
StepJob 中的一个独立阶段,通常对应一个读-处理-写操作序列。Step 是实际工作的执行单元,可配置事务边界、监听器、容错等。
ItemReader负责从数据源读取数据,每次读取一条记录(item)。实现应保证幂等性,以便支持重启。
ItemProcessor对读取的数据进行处理、转换或过滤。可返回 null 表示跳过该条记录(用于过滤)。
ItemWriter将处理后的数据写入目标位置,通常以块(chunk)为单位提交。写入操作应在事务内完成,确保一致性。
JobRepository持久化 Job 和 Step 的元数据(如执行状态、参数、开始/结束时间等)。依赖数据库存储,Spring Batch 自动管理其表结构。
JobLauncher用于启动 Job 的接口,负责创建 JobInstance 并执行。通常通过 Spring 自动注入使用。

1.4 批处理的执行流程

概念名称说明注意事项
JobInstance代表一个逻辑上的 Job 执行实例,由 Job 名称和 JobParameters 唯一标识。相同参数的 JobInstance 只能成功执行一次;不同参数可创建新实例。
JobParameters启动 Job 时传入的参数集合(如文件名、日期等),用于区分不同运行实例。参数是 JobInstance 唯一性的关键,建议包含时间戳或唯一标识。
JobExecution表示一次具体的 Job 执行过程,包含开始时间、结束时间、状态等信息。一个 JobInstance 可有多个 JobExecution(如失败后重启)。
StepExecution表示一个 Step 的一次执行,记录其读取、写入、跳过、错误等统计信息。每次 Step 执行都会生成一个 StepExecution,用于监控和恢复。
BatchStatus表示 Job 或 Step 的执行状态,如 STARTING、COMPLETED、FAILED 等。用于判断执行结果,决定是否重试或继续。
ExitStatusStep 的退出状态,可自定义(如 COMPLETED、FAILED、STOPPED 等)。可用于流程控制决策(如根据状态跳转到不同 Step)。

1.5 Spring Batch 与其他批处理框架的对比

框架名称特点与 Spring Batch 对比说明
Spring Batch基于 Spring 生态,配置灵活,支持数据库元数据管理,易于集成。优势:企业级支持、丰富的组件、良好的文档和社区。
Apache Camel集成框架,支持路由和中介模式,也可做批处理。Camel 更侧重 EIP(企业集成模式),而 Spring Batch 专注批处理核心逻辑。
Quartz调度框架,用于定时触发任务,但不提供读-处理-写模型。Quartz 可与 Spring Batch 结合使用(作为调度器),但本身不处理数据流。
JSR-352 (Java EE Batch)Java EE 标准批处理 API,功能类似 Spring Batch。Spring Batch 功能更丰富,社区活跃,而 JSR-352 实现较少且不够灵活。
Custom Scripts使用 Shell、Python 脚本实现批处理。脚本适合简单任务,缺乏监控、事务、重启等企业级特性。

第2章 环境搭建与第一个 Spring Batch 应用

2.1 搭建 Spring Boot + Spring Batch 项目(Maven/Gradle 配置)

构建工具配置说明注意事项
Maven使用 Spring Initializr 创建项目,选择 spring-boot-starter-batch 模块。确保父 POM 正确继承 spring-boot-starter-parent
Gradlebuild.gradle 中添加 implementation 'org.springframework.boot:spring-boot-starter-batch'使用正确的 Spring Boot 插件版本,避免依赖冲突。

2.2 引入必要的依赖

依赖名称用途说明注意事项
spring-boot-starter-batch核心依赖,包含 Spring Batch 所有功能模块。自动启用批处理配置,无需手动开启。
spring-boot-starter-jdbc 或 JPA提供数据访问能力,用于读写数据库。若使用数据库作为元数据存储,必须引入。
H2 / MySQL / Oracle Driver数据库驱动,根据实际使用的数据库选择。需在 application.properties 中正确配置连接信息。
spring-boot-starter-web(可选)若需提供 REST 接口启动 Job,可引入 Web 模块。非必需,但便于外部触发批处理任务。

2.3 配置数据源与批处理元数据表

配置项/概念说明注意事项
spring.datasource.*配置数据库连接信息(url、username、password、driver-class-name)。必须确保数据库可连接,否则 JobRepository 初始化失败。
spring.batch.jdbc.initialize-schema控制元数据表初始化行为:ALWAYSEMBEDDEDNEVERALWAYS:每次启动都建表(生产慎用);EMBEDDED:仅内存数据库时建表;NEVER:不建表。
元数据表(如 BATCH_JOB_INSTANCE)Spring Batch 使用约 19 张表记录 Job/Step 执行信息。表结构由框架自动管理,不应手动修改。
数据源共用 vs. 独立可使用业务数据源或独立的数据源存储批处理元数据。推荐使用独立数据源以避免资源竞争。

2.4 编写最简单的 Job 和 Step

概念/方法名语法用途说明代码示例注意事项
@EnableBatchProcessing@EnableBatchProcessing启用 Spring Batch 自动配置。@Configuration + @EnableBatchProcessing必须添加在配置类上,否则 Job 无法注册。
JobBuilderFactory@Autowired private JobBuilderFactory jobBuilderFactory;用于创建 Job 对象。Job job = jobBuilderFactory.get("simpleJob").start(step1).build();由 Spring 自动装配,无需手动 new。
StepBuilderFactory@Autowired private StepBuilderFactory stepBuilderFactory;用于创建 Step 对象。Step step = stepBuilderFactory.get("simpleStep").<String, String>chunk(10).reader(itemReader()).processor(itemProcessor()).writer(itemWriter()).build();chunk size 表示每批处理的记录数。
JobJob job = jobBuilderFactory.get("jobName").start(step).build();定义一个包含 Step 的 Job。见上例。Job 名称必须唯一。
StepStep step = stepBuilderFactory.get("stepName")...build();定义一个处理步骤。见上例。一个 Job 至少包含一个 Step。

2.5 运行并观察批处理执行日志

日志内容说明注意事项
Application StartedSpring Boot 应用启动完成,批处理自动执行(若未禁用)。默认情况下,Spring Batch 在应用启动时自动运行已定义的 Job。
Starting Job日志显示 Job 开始执行,包含 Job 名称和 JobParameters。可通过日志确认 Job 是否按预期启动。
Reading from ItemReader显示读取数据的日志(取决于日志级别)。可通过调试模式查看每条记录的读取情况。
Processed Item显示 ItemProcessor 处理结果(如有日志输出)。可用于验证数据转换逻辑是否正确。
Writing Items显示写入数据的数量和内容。确认写入是否成功,chunk 提交是否正常。
Finished JobJob 执行完成,输出最终状态(COMPLETED、FAILED 等)。成功状态应为 COMPLETED,失败则需排查异常堆栈。
异常堆栈(Exception Stack Trace)若执行失败,会输出详细异常信息。重点关注 Caused by 部分,定位根本原因。
Batch Metadata Logs记录 JobExecution、StepExecution 的持久化操作。表明元数据已写入数据库,支持后续重启。

第3章 Job 与 Step 的定义与控制

3.1 Job 的创建与配置(JobBuilderFactory)

方法名语法用途说明代码示例注意事项
getJobBuilder get(String jobName)获取一个 JobBuilder 实例,用于构建 Job。Job job = jobBuilderFactory.get("importUserJob").start(step1).build();jobName 必须唯一,否则会覆盖已有 Job。
startJobBuilder start(Step firstStep)指定 Job 的第一个 Step。jobBuilderFactory.get("job").start(step1).build();每个 Job 至少要有一个 Step。
nextJobBuilder next(Step nextStep)添加后续 Step,实现线性流程。jobBuilderFactory.get("job").start(step1).next(step2).build();Step 执行顺序由调用链决定。
flowJobBuilder flow(Step step)将 Step 包装为 Flow,用于复杂流程控制。jobBuilderFactory.get("job").start(flow1).build();需结合 end()on().to() 等方法实现条件跳转。
splitJobBuilder split(TaskExecutor taskExecutor)创建并行分支(Parallel Steps)。jobBuilderFactory.get("job").split(taskExecutor).add(flow1, flow2);需提供 TaskExecutor 实现并发执行。
incrementerJobBuilder incrementer(JobParametersIncrementer)设置参数增量器,避免 JobInstance 冲突。jobBuilderFactory.get("job").incrementer(new DailyJobIncrementer()).start(step);常用于每日执行的 Job,自动添加时间戳参数。
listenerJobBuilder listener(JobExecutionListener listener)注册 Job 级别的监听器。jobBuilderFactory.get("job").listener(jobListener).start(step);可多次调用添加多个监听器。
preventRestartJobBuilder preventRestart(boolean preventRestart)控制 Job 是否可重启。jobBuilderFactory.get("job").preventRestart(true).start(step);默认 false(可重启),设为 true 则失败后不能重启。
repositoryJobBuilder repository(JobRepository jobRepository)指定使用的 JobRepository。一般无需手动设置,由 Spring 自动注入。高级用法,用于自定义元数据存储。
buildJob build()构建并返回最终的 Job 实例。Job job = jobBuilderFactory.get("job").start(step).build();必须调用 build() 才能生成 Job 对象。

3.2 Step 的创建与配置(StepBuilderFactory)

方法名语法用途说明代码示例注意事项
getStepBuilder get(String stepName)获取 StepBuilder 实例。Step step = stepBuilderFactory.get("processDataStep").<String, String>chunk(10).reader(reader).writer(writer).build();stepName 在 Job 内应唯一。
chunkSimpleStepBuilder chunk(int chunkSize)配置基于块(chunk)的处理模式。stepBuilderFactory.get("step").<Integer, Integer>chunk(50).reader(itemReader).processor(itemProcessor).writer(itemWriter);chunkSize 表示每批处理多少条记录后提交事务。
taskletTaskletStepBuilder tasklet(Tasklet tasklet)配置 Tasklet 模式的 Step(非 chunk)。stepBuilderFactory.get("cleanUpStep").tasklet(cleanupTasklet).build();适用于无需读-处理-写循环的任务,如清理、通知等。
readerSimpleStepBuilder reader(ItemReader reader)设置数据读取器。见 chunk 示例。必须与 chunk 模式配合使用。
processorSimpleStepBuilder processor(ItemProcessor processor)设置数据处理器。见 chunk 示例。可选,若无需处理可省略。
writerSimpleStepBuilder writer(ItemWriter writer)设置数据写入器。见 chunk 示例。必须与 chunk 模式配合使用。
listenerStepBuilder listener(StepExecutionListener listener)注册 Step 级别的监听器。stepBuilderFactory.get("step").listener(stepListener).chunk(10)...可多次调用添加多个监听器。
transactionManagerStepBuilder transactionManager(PlatformTransactionManager tm)指定事务管理器。通常使用默认事务管理器,除非有特殊需求。在分布式事务场景下可能需要自定义。
allowStartIfCompleteStepBuilder allowStartIfComplete(boolean allow)是否允许 Step 已完成时仍可启动。stepBuilderFactory.get("step").allowStartIfComplete(true).chunk(10)...默认 false,设为 true 可重复执行已完成的 Step。
startLimitStepBuilder startLimit(int startLimit)设置 Step 最大启动次数。stepBuilderFactory.get("step").startLimit(3).chunk(10)...超过限制将抛出 StartLimitExceededException
buildStep build()构建并返回 Step 实例。Step step = stepBuilderFactory.get("step").chunk(10)...build();必须调用 build() 才能生成 Step 对象。

3.3 Step 的执行流程(read → process → write)

阶段说明注意事项
read由 ItemReader 从数据源逐条读取数据,直到返回 null 表示读取完成。每次调用 read() 返回一条记录,幂等性很重要,支持重启时从断点继续读取。
process将 read 返回的 item 传给 ItemProcessor 进行处理或过滤。若 processor 返回 null,该条记录将被跳过,不进入 write 阶段。
write将 process 返回的 item 集合(chunk size 条)一次性写入目标。write 在事务中执行,若失败则整个 chunk 回滚,可配置重试或跳过策略。
chunk 循环上述 read → process → write 以 chunk 为单位循环执行,直到数据读完。chunk size 影响性能和内存使用,需根据数据量和系统资源调优。
事务边界每个 chunk 的 write 操作在一个事务中提交。确保数据一致性,避免部分写入问题。

3.4 Job 的启动模式:启动新实例 vs. 继续上次执行

启动模式说明注意事项
启动新实例使用新的 JobParameters 创建新的 JobInstance,无论之前是否执行过。只要参数不同,即使 Job 名称相同,也会创建新实例。
继续上次执行使用与之前相同的 JobParameters,框架尝试恢复失败的 JobInstance。仅当 Job 未完成(如 FAILED 状态)时可继续;COMPLETED 的 Job 不能继续。
JobInstance 唯一性由 Job 名称 + JobParameters 共同决定。相同参数不能启动两个 JobInstance,除非设置 preventRestart=true
Restart 操作失败后使用相同参数重新运行 Job,从断点处恢复(如未完成的 Step)。需确保 ItemReader 支持重启(如记录当前读取位置)。
preventRestart若设置为 true,则 Job 一旦开始执行,即使失败也不能用相同参数重启。适用于不允许重复执行的敏感任务。

3.5 JobParameters 的使用与作用

方法/概念说明注意事项
JobParameters启动 Job 时传入的参数集合,类型包括 String、Long、Date、Double。参数值不可变,用于标识 JobInstance。
参数类型支持支持 STRING、LONG、DATE、DOUBLE 四种类型。DATE 类型格式为 yyyy/MMdd-hh:mm:ss
参数传递方式可通过命令行、REST API、编程方式(JobLauncher.run)传入。命令行示例:--fileName=input.csv --date=2025/09/28-10:00:00
JobParameters 作用1. 区分不同的 JobInstance
2. 作为数据源路径、过滤条件等动态输入
常用于指定输入文件名、处理日期等。
参数绑定可在 @Value 中使用 #{jobParameters['paramName']} 注入参数值。需配合 @StepScope@JobScope 使用。
JobParametersIncrementer自动为 JobParameters 添加增量值(如时间戳),避免参数重复。常见实现:RunIdIncrementerDailyJobTimeIncrementer

3.6 Job 的监听器(JobExecutionListener)

方法名语法用途说明代码示例注意事项
beforeJobvoid beforeJob(JobExecution jobExecution)Job 开始前执行的逻辑。见下方代码块可用于初始化资源、记录日志、发送通知等。
afterJobvoid afterJob(JobExecution jobExecution)Job 结束后执行的逻辑,无论成功或失败。见下方代码块可用于清理资源、发送执行结果邮件、更新状态等。
注册方式jobBuilderFactory.listener(listener)将监听器注册到 Job。jobBuilderFactory.get("job").listener(new MyJobListener()).start(step);可注册多个监听器,执行顺序按注册顺序。
jobExecution 参数包含 Job 的执行信息(状态、耗时、参数等)用于获取执行上下文。jobExecution.getStatus(), jobExecution.getEndTime()beforeJob 中 endTime 为 null。
public class MyJobListener implements JobExecutionListener {
    @Override
    public void beforeJob(JobExecution jobExecution) {
        System.out.println("Job 开始: " + jobExecution.getJobInstance().getJobName());
    }

    @Override
    public void afterJob(JobExecution jobExecution) {
        System.out.println("Job 结束: " + jobExecution.getStatus());
    }
}

3.7 Step 的监听器(StepExecutionListener)

方法名语法用途说明代码示例注意事项
beforeStepvoid beforeStep(StepExecution stepExecution)Step 开始前执行。见下方代码块可用于初始化 Step 相关资源。
afterStepExitStatus afterStep(StepExecution stepExecution)Step 结束后执行,可返回新的 ExitStatus。见下方代码块可用于修改 Step 的退出状态,影响后续流程决策。
注册方式stepBuilderFactory.listener(listener)将监听器注册到 Step。stepBuilderFactory.get("step").listener(new MyStepListener()).chunk(10)...可注册多个监听器。
stepExecution 参数包含 Step 的执行统计(读取数、写入数、跳过数等)用于监控和日志记录。stepExecution.getReadCount(), stepExecution.getWriteCount()beforeStep 中统计值为 0。
ExitStatus 返回值可返回 SUCCESS、FAILED、STOPPED 等状态影响 Job 流程走向。return ExitStatus.COMPLETED.and("customInfo");若返回 FAILED,Step 将标记为失败。
public class MyStepListener implements StepExecutionListener {
    @Override
    public void beforeStep(StepExecution stepExecution) {
        System.out.println("Step 开始: " + stepExecution.getStepName());
    }

    @Override
    public ExitStatus afterStep(StepExecution stepExecution) {
        System.out.println("Step 结束: " + stepExecution.getStatus());
        return ExitStatus.COMPLETED;
    }
}

3.8 Step 的执行结果控制(ExitStatus、BatchStatus)

状态类型说明注意事项
BatchStatus表示 Job 或 Step 的实际执行状态,由框架管理。包括:STARTING、STARTED、COMPLETED、FAILED、STOPPED、ABANDONED 等。
ExitStatus表示 Step 的退出状态,可自定义,用于流程决策。默认与 BatchStatus 一致,但可通过监听器或异常映射修改。
状态转换BatchStatus 是内部状态,ExitStatus 是对外暴露的状态。例如:Step 执行失败(BATCH_STATUS=FAILED),但 ExitStatus 可设为 COMPLETED(忽略错误)。
自定义 ExitStatus可通过 ExitStatus.and(String advice) 添加附加信息。如:ExitStatus.COMPLETED.and("file.processed")
异常映射使用 StepBuilder 的 faultTolerant().retry().skip() 配置异常处理策略。某些异常可配置为跳过或重试,不影响最终 BatchStatus。
流程决策可基于 ExitStatus 控制 Job 流程走向(如 on("COMPLETED").to(step2))。需结合 Job Flow 配置使用。

第4章 读取数据:ItemReader 详解

4.1 ItemReader 接口设计原理

概念/方法名说明注意事项
T read() throws Exception, UnexpectedInputException, ParseException, NonTransientResourceException;核心方法:每次调用返回一条记录,读取完毕返回 null。必须保证幂等性(Idempotent),即多次调用不会跳过或重复读取数据,以支持重启。
流式读取(Streaming)ItemReader 通常以流式方式逐条读取数据,避免一次性加载到内存。适用于大文件或大数据集,防止 OutOfMemoryError。
资源管理需实现 ResourceAwareItemReaderItemStream 接口来绑定数据源资源。如文件路径、数据库连接等,便于框架统一管理生命周期。
可重启性(Restartability)框架通过 ExecutionContext 保存读取位置(如行号、主键偏移量)。实现类需在重启时从断点继续读取,不能重新开始。
异常处理抛出 NonTransientResourceException 表示永久性错误,中断执行。Transient 异常可重试,NonTransient 则直接失败。
幂等性要求同一状态下调用 read() 多次应返回相同结果,直到状态改变。例如:第一次调用返回第1行,后续调用仍返回第1行,直到内部状态推进到第2行。

4.2 常用 ItemReader 实现类概述

实现类名称数据源类型特点说明适用场景
FlatFileItemReader文本文件、CSV、TSV支持按行读取,可配置分隔符、编码、跳过标题行等。日志文件、导出报表、数据导入
JdbcCursorItemReader关系型数据库使用 JDBC ResultSet 游标流式读取,内存占用低。大数据量查询,无需分页
JdbcPagingItemReader关系型数据库基于分页 SQL 查询(如 LIMIT/OFFSET 或数据库特定分页语法),适合大数据集。避免长事务,支持并行处理
StaxEventItemReaderXML 文件基于 SAX 的事件驱动解析,逐个读取 XML 元素,内存友好。结构化 XML 数据导入
RepositoryItemReaderJPA Repository调用 Spring Data JPA 的 PagingAndSortingRepository 分页读取实体。已有 JPA 实体和 Repository 的项目
ListItemReaderJava List从内存中的 List<T> 逐个读取数据,读完返回 null。单元测试、小数据集处理
MultiResourceItemReader多个文件包装其他 ItemReader,遍历多个资源(如多个 CSV 文件)。批量处理同格式的多个文件
ItemReaderAdapter任意服务方法将普通服务方法适配为 ItemReader,通过反射调用远程或本地接口。调用遗留系统、Web Service 获取数据

4.3 FlatFileItemReader:读取文本/CSV 文件

配置属性/方法说明注意事项
setResource(Resource resource)设置输入文件资源。支持 ClassPathResource、FileSystemResource 等。
setLineMapper(LineMapper<T> lineMapper)定义如何将一行文本映射为对象,需配合 DefaultLineMapper 使用。核心配置项。
setLinesToSkip(int linesToSkip)跳过前 N 行(如标题行)。避免将表头当作数据处理。
setEncoding(String encoding)设置文件编码(如 UTF-8)。中文文件建议显式设置。
setStrict(boolean strict)是否严格模式:文件不存在时是否抛异常。设为 false 可容忍文件缺失。

DefaultLineMapper + DelimitedLineTokenizer CSV 解析示例:

DefaultLineMapper<User> lineMapper = new DefaultLineMapper<>();
DelimitedLineTokenizer tokenizer = new DelimitedLineTokenizer(",");
tokenizer.setNames("id", "name", "email");
lineMapper.setLineTokenizer(tokenizer);
lineMapper.setFieldSetMapper(fieldSet ->
    new User(fieldSet.readLong("id"), fieldSet.readString("name")));
reader.setLineMapper(lineMapper);
  • FixedLengthTokenizer:用于定长文件解析(银行、金融等传统系统输出的固定宽度文件),需定义每个字段的起始和结束位置。

4.4 JdbcCursorItemReader:基于 JDBC 游标的数据库读取

配置属性/方法说明注意事项
setDataSource(DataSource dataSource)设置数据源。通常由 Spring 容器注入。
setSql(String sql)设置查询 SQL。不支持 LIMIT/OFFSET,全量查询。
setRowMapper(RowMapper<T> rowMapper)将 ResultSet 行映射为对象。类似 Spring JDBC 的 RowMapper。
setPreparedStatementSetter(PreparedStatementSetter pss)设置 SQL 参数。用于设置 WHERE 条件参数。
游标模式使用数据库游标流式读取,不缓存全部结果集。内存占用低,依赖数据库驱动支持游标,且事务需保持打开状态。
性能特点高吞吐量,但占用数据库连接时间较长。适合单线程处理大批量数据,不适合高并发场景。

4.5 JdbcPagingItemReader:分页读取数据库

配置属性/方法说明注意事项
setPageSize(int pageSize)每页记录数。建议与 Step 的 chunk size 一致。
setSelectClause(String selectClause)SELECT 子句。不包含 SELECT 关键字也可。
setFromClause(String fromClause)FROM 子句。分离 SQL 片段便于拼接。
setWhereClause(String whereClause)WHERE 条件。支持命名参数(:param)。
setSortKeys(Map<String, Order> sortKeys)排序列及顺序,用于分页定位。必须配置,否则无法保证分页一致性。
setParameterValues(Map<String, Object> parameterValues)设置查询参数。支持动态参数绑定。
分页策略不同数据库使用不同分页语法,框架自动适配。需正确配置 DatabaseType。
优点事务短,连接释放快,支持并行 Step(split)。排序列必须唯一且稳定,避免漏读或重复读。

4.6 StaxEventItemReader:读取 XML 文件

配置属性/方法说明注意事项
setResource(Resource resource)设置 XML 文件资源。文件格式需符合配置的结构。
setFragmentRootElementName(String fragmentRootElementName)指定要读取的元素标签名(碎片根元素)。每次读取一个 <user>...</user> 片段。
setUnmarshaller(Unmarshaller unmarshaller)使用 JAXB 或其他工具将 XML 片段反序列化为对象。需配置绑定类,字段需有 @XmlElement 等注解。
XML 结构要求文件应为 <root><user>...</user><user>...</user></root> 形式。不支持复杂嵌套结构的逐条读取,每个 fragment 应独立可解析。
内存效率基于事件驱动(SAX),不加载整个文档,内存友好。适合大 XML 文件处理,比 DOM 解析更高效。

4.7 RepositoryItemReader:通过 Spring Data JPA 读取

配置属性/方法说明注意事项
setRepository(PagingAndSortingRepository repository)设置 JPA Repository。必须是 PagingAndSortingRepository 或其子接口。
setQueryMethod(String queryMethod)指定 Repository 中的查询方法名。方法需支持 Pageable 参数。
setArguments(List<?> arguments)传递给查询方法的参数列表。顺序必须与方法参数一致。
setSort(Map<String, Sort.Direction> sort)排序列及方向。必须配置,用于分页定位。
setPageSize(int pageSize)每页大小。与 JPA 分页一致。
优点与现有 JPA 代码无缝集成,复用 Repository 逻辑,减少重复 SQL 编写。需确保 Repository 方法支持分页。
依赖必须引入 spring-boot-starter-data-jpa并配置好 EntityManager 和事务管理器,否则无法工作。

4.8 自定义 ItemReader 的实现

实现场景说明注意事项
从远程 API 读取调用 RESTful 接口分页获取数据。见下方代码块
读取 NoSQL 数据如 MongoDB、Redis 中的数据流。注意 NoSQL 的分页方式(如 MongoDB 的 skip/limit 或游标)。
实现 ItemStream 接口支持重启时保存和恢复状态(如当前页码、文件偏移量)。必须在 update 中保存状态,在 open 中恢复。
异常处理抛出 NonTransientResourceException 表示不可恢复错误。避免捕获异常后返回 null 导致误判为”读取完成”。
资源释放实现 close() 方法释放连接、流等资源。防止资源泄漏。
public class ApiItemReader implements ItemReader<User>, ItemStream {
    private Iterator<User> currentIterator;
    private int currentPage = 0;
    private final int pageSize = 100;

    @Override
    public User read() throws Exception {
        if (currentIterator == null || !currentIterator.hasNext()) {
            List<User> page = fetchPage(currentPage++, pageSize);
            if (page.isEmpty()) {
                return null; // 读取完毕
            }
            currentIterator = page.iterator();
        }
        return currentIterator.next();
    }

    private List<User> fetchPage(int page, int size) {
        // 调用远程 API 获取数据
    }

    @Override
    public void open(ExecutionContext executionContext) {
        currentPage = executionContext.getInt("currentPage", 0);
    }

    @Override
    public void update(ExecutionContext executionContext) {
        executionContext.putInt("currentPage", currentPage);
    }

    @Override
    public void close() {
        // 释放资源
    }
}

第5章 处理数据:ItemProcessor 详解

5.1 ItemProcessor 接口作用与设计模式

概念/方法名说明注意事项
T process(T item) throws Exception;核心方法:接收一个输入对象,返回处理后的对象或 null。方法必须是幂等的,以支持重试和重启。
作用在 ItemReader 和 ItemWriter 之间对数据进行转换、过滤或增强。是 Spring Batch 中实现业务逻辑的核心组件之一。
设计模式责任链模式(Chain of Responsibility):可串联多个处理器。每个处理器只关注单一职责,便于测试和维护。
返回值含义返回对象:继续传递给下一个组件;返回 null:过滤掉该条数据,不写入利用 null 实现条件过滤是常见做法。
无状态要求默认应设计为无状态,避免在实例变量中保存上下文。若需状态,应结合 ItemStream 接口管理执行上下文。
线程安全在多线程 Step 中,ItemProcessor 必须是线程安全的。避免使用可变实例变量,推荐使用不可变对象或局部变量。

5.2 实现简单的数据转换与过滤

功能类型说明代码示例
数据转换将一种类型转换为另一种类型,或修改字段值。见下方数据转换示例
条件过滤根据业务规则决定是否保留该条数据(返回 null 即过滤)。见下方条件过滤示例
字段增强添加计算字段、默认值或外部信息(如调用服务)。见下方字段增强示例

数据转换示例:

public class UserToDtoProcessor implements ItemProcessor<User, UserDto> {
    @Override
    public UserDto process(User user) {
        UserDto dto = new UserDto();
        dto.setId(user.getId());
        dto.setFullName(user.getFirstName() + " " + user.getLastName());
        dto.setEmail(user.getEmail().toLowerCase());
        return dto;
    }
}

条件过滤示例:

public class ActiveUserFilter implements ItemProcessor<User, User> {
    @Override
    public User process(User user) {
        return "ACTIVE".equals(user.getStatus()) ? user : null;
    }
}

字段增强示例:

public class EnrichUserProcessor implements ItemProcessor<User, User> {
    @Autowired
    private AddressService addressService;

    @Override
    public User process(User user) {
        Address address = addressService.getByUserId(user.getId());
        user.setAddress(address);
        return user;
    }
}

⚠️ 注意:@Autowired 在无状态处理器中可用,但需确保线程安全。

5.3 ValidatingItemProcessor:数据校验处理器

配置/使用方式说明注意事项
ValidatingItemProcessor包装一个 Validator<T>,验证失败时抛出 ValidationException必须设置 setFilter(true) 才能实现”静默过滤”,否则会中断批处理。
自定义 Validator实现 org.springframework.validation.Validator 接口。见下方自定义 Validator 示例
使用 Bean Validation结合 javax.validation 注解(如 @NotNull, @Email)。需引入 spring-boot-starter-validation

ValidatingItemProcessor 使用示例:

ValidatingItemProcessor<User> validator =
    new ValidatingItemProcessor<>(new UserValidator());
validator.setFilter(true); // 验证失败时返回 null,而非抛异常

自定义 Validator 示例:

public class UserValidator implements Validator {
    @Override
    public boolean supports(Class<?> clazz) {
        return User.class.equals(clazz);
    }

    @Override
    public void validate(Object target, Errors errors) {
        User user = (User) target;
        if (user.getEmail() == null || !user.getEmail().contains("@")) {
            errors.rejectValue("email", "invalid", "Email is invalid");
        }
    }
}

Bean Validation 示例:

public class User {
    @NotNull
    private String name;
    @Email
    private String email;
    // getters/setters
}

// 使用:
ValidatingItemProcessor<User> processor =
    new ValidatingItemProcessor<>(new LocalValidatorFactoryBean());
processor.setFilter(true);

5.4 CompositeItemProcessor:组合多个处理器

特性说明注意事项
功能将多个 ItemProcessor 串联执行,形成处理链。处理器顺序很重要,前一个的输出是后一个的输入。
类型安全支持泛型链式转换(如 User → UserDto → EnrichedDto)。编译期检查类型匹配,避免运行时错误;若类型不匹配会抛 ClassCastException
配置方式可通过 Java Config 或 XML 配置。Java Config 更清晰,推荐使用。
异常传播任一处理器抛出异常都会中断整个链。可结合容错机制(skip/retry)提升容错能力;过滤(返回 null)会在链中传递,最终导致不写入。
List<ItemProcessor<?, ?>> delegates = Arrays.asList(
    new UserToDtoProcessor(),
    new ActiveUserFilter(),
    new EnrichUserProcessor()
);

CompositeItemProcessor<User, UserDto> processor = new CompositeItemProcessor<>();
processor.setDelegates(delegates);

5.5 容错处理(Skip / Retry)

⚠️ 注意:Spring Batch 原生没有 FaultTolerantItemProcessor 类。容错(如跳过、重试)是 Step 级别 的配置,作用于整个 ItemProcessor 或 ItemWriter。

容错机制配置方式说明
跳过(Skip).faultTolerant().skip(ValidationException.class).skipLimit(10)当 ItemProcessor 抛出指定异常(如校验失败),跳过该条记录并计入 skippedCount。适用于可容忍少量错误数据的场景。
重试(Retry).faultTolerant().retry(DeadlockLoserDataAccessException.class).retryLimit(3)对于临时性异常(如死锁、网络抖动),自动重试处理该条记录。需确保 process() 方法幂等。
自定义 SkipPolicy实现 SkipPolicy 接口,自定义跳过逻辑。比简单 skip(Class) 更灵活,可结合异常类型和次数判断。

Step 级别容错配置示例:

@Bean
public Step myStep(ItemReader<User> reader,
                   ItemProcessor<User, UserDto> processor,
                   ItemWriter<UserDto> writer) {
    return stepBuilderFactory.get("myStep")
        .<User, UserDto>chunk(10)
        .reader(reader)
        .processor(processor)
        .writer(writer)
        .faultTolerant()
        .skip(ValidationException.class)
        .skipLimit(10)
        .retry(DeadlockLoserDataAccessException.class)
        .retryLimit(3)
        .build();
}

自定义 SkipPolicy 示例:

public class CustomSkipPolicy implements SkipPolicy {
    @Override
    public boolean shouldSkip(Throwable t, int skipCount) {
        if (t instanceof ValidationException && skipCount < 5) {
            return true; // 最多跳过 5 次校验异常
        }
        return false;
    }
}

5.6 自定义 ItemProcessor 实现

实现场景说明注意事项
调用外部服务如调用 REST API 补全数据。注意异常处理、超时设置、线程安全。
批量预加载缓存在处理前加载字典数据,避免每条记录都查数据库。使用 @BeforeStep 初始化缓存,提升性能。
实现 ItemStream支持重启时恢复状态(如已处理条数、缓存快照)。实现 open(), update(), close() 方法,将状态存入 ExecutionContext。
日志与监控记录处理日志、埋点、指标统计。避免频繁 I/O 影响性能,可异步记录。

调用外部服务示例:

public class ApiEnrichingProcessor implements ItemProcessor<Order, Order> {
    private final RestTemplate restTemplate;

    public ApiEnrichingProcessor(RestTemplate restTemplate) {
        this.restTemplate = restTemplate;
    }

    @Override
    public Order process(Order order) {
        try {
            Customer customer = restTemplate.getForObject(
                "/api/customers/{id}", Customer.class, order.getCustomerId());
            order.setCustomer(customer);
            return order;
        } catch (HttpClientErrorException.NotFound e) {
            return null; // 客户不存在则过滤
        }
    }
}

批量预加载缓存示例:

public class CachedRoleProcessor implements ItemProcessor<User, User> {
    private Map<Long, String> roleCache;

    @BeforeStep
    public void beforeStep(StepExecution stepExecution) {
        roleCache = roleRepository.findAll().stream()
            .collect(Collectors.toMap(Role::getId, Role::getName));
    }

    @Override
    public User process(User user) {
        user.setRoleName(roleCache.get(user.getRoleId()));
        return user;
    }
}

最佳实践总结:

  • 优先使用组合模式(CompositeItemProcessor)而非大而全的处理器。
  • 过滤逻辑推荐返回 null,配合 skip 策略实现优雅容错。
  • 复杂逻辑可拆分为多个小处理器,便于单元测试。
  • 外部依赖(如 API、DB)应添加超时、降级、重试机制。

第6章 写入数据:ItemWriter 详解

6.1 ItemWriter 接口设计原理

概念/方法名说明注意事项
void write(List<? extends T> items) throws Exception;核心方法:接收一个项目列表并写入目标。方法名虽为 write,但参数是 List,体现批处理特性。
批量写入(Chunk-oriented)ItemWriter 一次性处理一个 chunk(块)的数据,而非单条记录。chunk 大小由 Step 的 chunkSize 决定,如 chunkSize=10,则每次 write() 接收最多10条数据。
事务边界每次 write() 调用通常在一个事务内完成。若写入失败,整个 chunk 回滚,保证数据一致性。
无返回值与 ItemProcessor 不同,ItemWriter 不返回值。写入成功即完成,失败则抛异常触发容错机制。
资源管理应实现 ItemStream 接口以管理资源(如文件流、数据库连接)。框架在 Step 开始时调用 open(),结束时调用 close()
幂等性要求write() 应设计为幂等,支持重启和重试。例如:使用 MERGE/UPSERT 而非 INSERT 避免重复。

6.2 常用 ItemWriter 实现类概述

实现类名称数据目标特点说明适用场景
FlatFileItemWriter文本文件、CSV、TSV将对象列表格式化为文本行写入文件。支持头部、尾部、分隔符等。数据导出、生成报表、文件交换
JdbcBatchItemWriter关系型数据库使用 JDBC PreparedStatement 批量插入/更新,性能高。大数据量写入数据库,ETL 场景
JpaItemWriterJPA 实体管理通过 EntityManager 批量持久化实体,支持级联。已使用 JPA 的项目,需复杂对象图保存
CompositeItemWriter多个目标包装多个 ItemWriter,将同一数据写入多个目的地。数据同步、备份、多格式输出
ClassifierCompositeItemWriter分类写入根据分类器(Classifier)将不同数据路由到不同 ItemWriter。条件分发,如按地区写入不同表/文件
ItemWriterAdapter任意服务方法将普通服务方法适配为 ItemWriter,通过反射调用。调用遗留系统、外部 API 写入数据
MongoItemWriterMongoDB写入文档到 MongoDB 集合。NoSQL 场景,JSON 数据存储
KafkaItemWriterKafka Topic将数据发送到 Kafka 消息队列。实时数据管道、事件驱动架构

6.3 FlatFileItemWriter:写入文本/CSV 文件

配置属性/方法说明注意事项
setResource(Resource resource)设置输出文件资源。支持 FileSystemResource 等。
setLineAggregator(LineAggregator<T> aggregator)定义如何将对象转换为一行文本。必须配置,核心组件。
setShouldDeleteIfExists(boolean shouldDeleteIfExists)文件已存在时是否删除。避免追加旧数据,确保结果纯净。
setAppendAllowed(boolean appendAllowed)是否允许追加写入。通常设为 false,避免数据混乱。
setHeaderCallback(FlatFileHeaderCallback)写入文件头部(如列名)。生成标准 CSV 文件。
setFooterCallback(FlatFileFooterCallback)写入文件尾部(如统计信息)。添加元数据、时间戳等。
setEncoding(String encoding)设置文件编码。中文内容必须设置,避免乱码。

DelimitedLineAggregator 配置示例:

DelimitedLineAggregator<User> lineAggregator = new DelimitedLineAggregator<>();
lineAggregator.setDelimiter(",");
lineAggregator.setFieldExtractor(item -> new String[]{
    String.valueOf(item.getId()),
    item.getName(),
    item.getEmail()
});
writer.setLineAggregator(lineAggregator);

Header / Footer 示例:

writer.setHeaderCallback(writer -> writer.write("id,name,email"));
writer.setFooterCallback(writer -> writer.write("END," + new Date()));

6.4 JdbcBatchItemWriter:批量写入数据库(JDBC)

配置属性/方法说明注意事项
setDataSource(DataSource dataSource)设置数据源。从 Spring 容器注入。
setSql(String sql)设置 INSERT 或 UPDATE SQL。使用 ? 占位符,不可用命名参数。
setItemPreparedStatementSetter(ItemPreparedStatementSetter<T>)设置参数值。或使用 BeanPropertyItemSqlParameterSourceProvider(见下)。
setItemSqlParameterSourceProvider(ItemSqlParameterSourceProvider<T>)基于对象属性自动绑定参数(推荐)。简化代码,要求字段名与 SQL 占位符对应。
setAssertUpdates(boolean assertUpdates)是否验证每条 SQL 的影响行数。设为 true 时,若某条 INSERT 影响0行会抛异常。
批量提交使用 JdbcTemplate.batchUpdate() 实现。框架自动处理,性能远高于单条 INSERT。
事务管理在 Step 的事务内执行,失败回滚。确保 DataSource 支持事务,避免在 ItemWriter 内开启新事务。

配置示例:

writer.setSql("INSERT INTO user_export (id, name, email) VALUES (?, ?, ?)");
writer.setItemPreparedStatementSetter((user, ps) -> {
    ps.setLong(1, user.getId());
    ps.setString(2, user.getName());
    ps.setString(3, user.getEmail());
});

// 或使用 BeanPropertyItemSqlParameterSourceProvider(推荐)
writer.setItemSqlParameterSourceProvider(
    new BeanPropertyItemSqlParameterSourceProvider<>());

6.5 JpaItemWriter:通过 JPA 写入实体

配置属性/方法说明注意事项
setEntityManagerFactory(EntityManagerFactory emf)设置 EMF。由 Spring JPA 自动配置。
setUsePersist(boolean usePersist)使用 persist() 还是 merge()persist 用于新实体,merge 用于更新。
setClearPersistenceContext(boolean clear)每次写入后是否清除 EntityManager 一级缓存。强烈建议设为 true,防止内存溢出(OutOfMemoryError)。
批量处理write(List<T>) 中遍历调用 persist()flush()框架自动处理。
性能优化配置 hibernate.jdbc.batch_sizehibernate.order_inserts启用 JDBC 批量插入,显著提升性能。
适用场景需要 JPA 级联保存、监听器(@PrePersist)、复杂对象图。简单插入推荐 JdbcBatchItemWriter,性能更高。

JPA 批量优化配置(application.yml):

spring:
  jpa:
    properties:
      hibernate:
        jdbc:
          batch_size: 50
        order_inserts: true

6.6 CompositeItemWriter:组合多个写入器

特性说明注意事项
功能将同一份数据同时写入多个目标。所有 write() 调用会依次执行。
事务行为所有写入在一个事务内,任一失败则全部回滚。保证数据一致性;若需独立事务,应使用 Step 的 split 并行执行。
性能影响写入时间 = 所有 ItemWriter 写入时间之和。可能成为瓶颈,考虑异步写入或选择性组合。
典型用例数据导出到文件 + 写入数据库;写入主库 + 发送到消息队列实现数据冗余、审计、解耦。
List<ItemWriter<UserDto>> writers = Arrays.asList(fileWriter, dbWriter);

CompositeItemWriter<UserDto> compositeWriter = new CompositeItemWriter<>();
compositeWriter.setDelegates(writers);

6.7 自定义 ItemWriter 实现

实现场景说明注意事项
写入 NoSQL如 Redis、Elasticsearch。注意序列化格式、连接池、错误处理。
调用外部 API如调用 RESTful 服务批量创建资源。建议实现批量 API(如 /api/orders/batch)以减少调用次数。
实现 ItemStream管理资源(如文件流、网络连接)。必须正确实现 open()close(),防止资源泄漏。
容错处理write() 内部处理部分失败。捕获异常,记录失败项,继续处理其余数据。

写入 Redis 示例:

public class RedisItemWriter implements ItemWriter<User> {
    private final StringRedisTemplate redisTemplate;

    public RedisItemWriter(StringRedisTemplate redisTemplate) {
        this.redisTemplate = redisTemplate;
    }

    @Override
    public void write(List<? extends User> users) {
        users.forEach(user ->
            redisTemplate.opsForValue().set(
                "user:" + user.getId(),
                JsonUtils.toJson(user)
            )
        );
    }
}

调用外部 API 示例:

public class ApiItemWriter implements ItemWriter<Order> {
    private final RestTemplate restTemplate;

    @Override
    public void write(List<? extends Order> orders) {
        orders.forEach(order ->
            restTemplate.postForEntity("/api/orders", order, Void.class)
        );
    }
}

实现 ItemStream 示例:

public class StreamingFileWriter implements ItemWriter<String>, ItemStream {
    private BufferedWriter writer;
    private Resource resource;

    @Override
    public void open(ExecutionContext executionContext) {
        try {
            File file = resource.getFile();
            writer = Files.newBufferedWriter(file.toPath(), StandardCharsets.UTF_8);
        } catch (IOException e) {
            throw new ItemStreamException("Cannot open file", e);
        }
    }

    @Override
    public void write(List<? extends String> lines) throws Exception {
        for (String line : lines) {
            writer.write(line);
            writer.newLine();
        }
    }

    @Override
    public void close() {
        if (writer != null) {
            try { writer.close(); }
            catch (IOException e) {
                throw new ItemStreamException("Cannot close file", e);
            }
        }
    }
}

最佳实践总结:

  • 优先选择 JdbcBatchItemWriter 进行数据库写入,性能最优。
  • JpaItemWriter 仅在需要 JPA 特性(如级联、监听器)时使用,并务必启用 clearPersistenceContext 和批量插入。
  • FlatFileItemWriter 注意编码和文件覆盖策略。
  • CompositeItemWriter 用于强一致性场景;若目标系统独立,建议用 Step Split 并行处理。
  • 自定义 ItemWriter 必须实现 ItemStream 管理资源,并保证幂等性。

第7章 批处理的容错与重启机制

7.1 重试机制(Retry)配置与使用

配置项/概念说明注意事项
.faultTolerant()启用容错模式,是使用 retry 和 skip 的前提。不启用则无法使用重试/跳过。
.retry(Class<? extends Throwable> exceptionType)指定哪些异常触发重试。只对指定异常重试,其他异常仍会中断 Step。
.retryLimit(int limit)设置最大重试次数。包括首次执行,总尝试次数 = retryLimit + 1。例如 limit=3,最多尝试4次。
重试时机在同一条记录上重复执行 ItemProcessor.process()ItemWriter.write()对于 ItemWriter,重试的是整个 chunk,而非单条记录。
幂等性要求process()write() 必须幂等,避免重复操作产生副作用(如重复扣款)。使用数据库乐观锁、唯一约束、状态机等保证;非幂等操作严禁使用重试!
自定义 RetryCallback实现复杂重试逻辑(如退避策略)。可在 ItemProcessor 或 ItemWriter 内部使用 RetryTemplate。

Step 级别重试配置示例:

return stepBuilderFactory.get("retryStep")
    .<User, UserDto>chunk(10)
    .reader(itemReader())
    .processor(itemProcessor())
    .writer(itemWriter())
    .faultTolerant() // 必须先启用容错
    .retry(DeadlockLoserDataAccessException.class)
    .retry(HttpServerErrorException.class)
    .retryLimit(3)
    .build();

自定义 RetryTemplate 退避策略示例:

RetryTemplate retryTemplate = new RetryTemplate();
retryTemplate.setRetryPolicy(new SimpleRetryPolicy(3));
retryTemplate.setBackOffPolicy(new ExponentialBackOffPolicy());

retryTemplate.execute(context -> {
    // 调用可能失败的服务
    externalService.call();
    return null;
});

7.2 跳过记录(Skip)策略配置

配置项/概念说明注意事项
.skip(Class<? extends Throwable> exceptionType)指定哪些异常应被跳过(即忽略该条记录)。异常发生时,当前记录被标记为”skipped”并继续处理下一条。
.skipLimit(int limit)允许跳过的最大记录数。达到限制后再次遇到可跳过异常,将抛出异常终止 Step。
.noSkip(Class<? extends Throwable>)明确排除某些异常不被跳过。即使该异常类型在 skip() 列表中,也会中断流程。
跳过行为ItemProcessor 返回 null:框架视为”过滤”,计入 filterCount,不计入 skipCount。理解 filter vs skip 的统计区别。
访问跳过计数在监听器中获取 StepExecution 的 getSkipCount()用于日志、告警、后续处理。

在监听器中获取跳过计数:

@AfterStep
public ExitStatus afterStep(StepExecution stepExecution) {
    long skipped = stepExecution.getSkipCount();
    if (skipped > 0) {
        log.warn("Skipped {} records", skipped);
    }
    return stepExecution.getExitStatus();
}

7.3 重试与跳过的异常分类

异常类型推荐处理方式常见异常举例说明
临时性异常 (Transient)✅ 重试 (retry)DeadlockLoserDataAccessException, HttpTimeoutException, ConnectionLossException由瞬时故障引起,重试可能成功。
数据校验异常 (Validation)✅ 跳过 (skip)ValidationException, ConstraintViolationException, NumberFormatException数据本身有问题,重试无意义,应跳过坏数据。
资源缺失异常 (Resource Not Found)⚠️ 根据场景选择FileNotFoundException, NoSuchElementException若文件偶尔缺失可重试;若记录不存在,通常跳过。
致命异常 (Fatal)❌ 不处理(中断)OutOfMemoryError, StackOverflowError, BeanCreationException应立即停止,需人工干预。
业务规则异常 (Business Rule)⚠️ 根据策略InsufficientFundsException, OrderAlreadyShippedException可能需要跳过或作为致命错误,取决于业务。
权限/配置异常❌ 不处理AccessDeniedException, IllegalConfigurationException配置错误,重试无效。

📊 组合策略示例:

.faultTolerant()
    .retry(DeadlockLoserDataAccessException.class)
    .retry(HttpServerErrorException.InternalServerError.class)
    .retryLimit(3)
    .skip(ValidationException.class)
    .skip(NumberFormatException.class)
    .skipLimit(50) // 最多容忍50条脏数据
    .noSkip(SystemUnavailableException.class) // 此异常必须中断

7.4 重启 Job:JobInstance 与 JobParameters 的关系

概念说明关键点
JobInstance一个作业定义(Job)的一次逻辑执行实例。由 Job 名称 + JobParameters 唯一确定;同一个 JobInstance 可以多次运行(重启),形成多个 JobExecution。
JobExecution一次具体的物理执行过程,属于一个 JobInstance。每次启动(无论是否重启)都会创建新的 JobExecution;包含开始时间、结束时间、状态(COMPLETED, FAILED, STOPPED)、退出码等。
JobParameters启动 Job 时传入的一组参数,用于区分不同的 JobInstance。核心作用:决定是否创建新实例 or 重启旧实例;类型:String, Long, Date, Double;可设置 identifying=false 的参数不影响实例唯一性。
重启判定逻辑当使用完全相同的 JobParameters 再次启动 Job 时框架查询数据库,查找已存在的 JobInstance。如果存在且上次执行未完成(状态为 FAILED 或 STOPPED),则重启该实例,从断点处继续;如果不存在或上次已 COMPLETED,则创建新实例。
断点续做Spring Batch 将每个 Step 的执行状态存储在 ExecutionContext 中。重启时,ItemReader 从上次位置继续读取,实现”断点续传”。

7.5 防止重复执行的参数控制策略

策略说明优缺点
使用时间戳/序列号确保每次运行的 JobParameters 唯一。✅ 优点:简单可靠,绝对防止重复。❌ 缺点:无法重启,每次都是新实例。
使用业务日期以处理的数据日期作为参数。✅ 优点:语义清晰,支持按天重试。✅ 可重启:若当天任务失败,用相同参数重启。
identifying=false 参数添加不影响实例唯一性的参数用于控制。✅ 优点:可在不创建新实例的前提下传递额外信息。❌ 缺点:修改 identifying=false 参数不会触发新实例。
数据库状态检查在 Job 启动前,通过自定义逻辑检查是否允许运行。✅ 优点:最灵活。❌ 缺点:增加复杂性。
Spring Batch Admin / 调度器控制使用外部调度工具并配置 concurrent="false"✅ 优点:职责分离。✅ 推荐生产环境使用。
// 时间戳方式
JobParameters params = new JobParametersBuilder()
    .addLong("run.id", System.currentTimeMillis())
    .addDate("date", new Date())
    .toJobParameters();

// 业务日期方式
JobParameters params = new JobParametersBuilder()
    .addDate("batch.date", targetDate) // 如 "2025-09-28"
    .toJobParameters();

// identifying=false 参数
JobParameters params = new JobParametersBuilder()
    .addDate("targetDate", targetDate) // identifying=true (默认)
    .addString("mode", "test", false) // identifying=false
    .toJobParameters();

最佳实践总结:

  • 明确需求:需要重启?还是绝对防重?
  • 推荐方案:对于周期性批处理(如每日报表),使用 业务日期 作为 JobParameter。
  • 关键原则:JobParameters 是控制重启与防重的核心,设计时务必谨慎。
  • 监控:利用 JobRepository 查询 JobInstance 和 JobExecution 状态,确保按预期运行。

第8章 批处理的性能优化与分片处理

8.1 提高 Step 的处理性能:chunk size 调优

概念说明调优建议/注意事项
Chunk Size每次 ItemWriter.write() 接收的项目数量,是批处理的核心性能参数。默认值通常为 10,需根据场景调整;chunkSize 影响事务边界:一个 chunk = 一个事务。
过小的 chunkSize(1-5)❌ 缺点:事务开销大(频繁提交/回滚)、I/O 效率低、总体吞吐量低。仅适用于对一致性要求极高、数据量极小的场景。
过大的 chunkSize(> 1000)❌ 缺点:单次事务过长、锁竞争激烈、内存压力大、失败重试成本高。风险高,不推荐盲目增大。
合理范围✅ 建议:数据库写入 50-500;文件写入 100-1000;外部 API 调用根据 API 限制(如批量 10-100)。需结合 8.5 节的分页优化。
调优方法基准测试 → 监控指标(readCount, writeCount, commitCount, GC 频率)→ 逐步调整。选择在稳定吞吐量和可接受失败成本之间的最优值。

8.2 并行 Step(Parallel Steps)配置

概念说明注意事项
Parallel Steps将独立的 Step 在同一个 Job 中并行执行,缩短总运行时间。split() 内的 Step 必须逻辑独立,无数据依赖。
执行器 (TaskExecutor)负责并发执行 Step,推荐使用 ThreadPoolTaskExecutor使用 SimpleAsyncTaskExecutor 或自定义 TaskExecutor。
适用场景多个数据源的独立 ETL 任务、生成不同报表、并行调用多个微服务。是 Job 级别的并行化。
@Bean
public Job parallelJob() {
    return jobBuilderFactory.get("parallelJob")
        .start(parallelFlow())
        .next(finalStep()) // 所有并行步骤完成后执行
        .end()
        .build();
}

private Flow parallelFlow() {
    return new FlowBuilder<SimpleFlow>("parallelFlow")
        .split(taskExecutor()) // 使用线程池执行器
        .add(
            step1(), // 可独立运行
            step2()  // 可独立运行
        )
        .build();
}

private TaskExecutor taskExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(4);
    executor.setMaxPoolSize(4);
    executor.setQueueCapacity(10);
    executor.setThreadNamePrefix("parallel-step-");
    executor.initialize();
    return executor;
}

8.3 分片 Step(Partitioning Step)原理与实现

概念说明注意事项
Partitioning Step将一个大 Step 拆分为多个分片(Shard),并行处理不同数据子集。解决单 Step 瓶颈,实现数据级并行。
核心组件Partitioner、Step (Worker)、TaskExecutor、PartitionHandler分片间必须无数据依赖。
数据划分方式基于主键范围(推荐)、基于列表、基于时间范围、基于哈希。选择能均匀分布数据的划分方式。
优势✅ 可处理远超单机内存的数据量、充分利用多核 CPU 和 I/O 带宽、可扩展性强。复杂度高,需仔细设计 Partitioner。
// 1. 定义 Partitioner
@Bean
public Partitioner partitioner() {
    return gridSize -> {
        Map<String, ExecutionContext> result = new HashMap<>();
        int range = 10000 / gridSize; // 假设总数据 10000 条
        for (int i = 0; i < gridSize; i++) {
            ExecutionContext context = new ExecutionContext();
            context.putInt("minValue", i * range);
            context.putInt("maxValue", (i + 1) * range);
            result.put("partition" + i, context);
        }
        return result;
    };
}

// 2. 定义 Worker Step(可复用普通 Step)
@Bean
public Step workerStep() {
    return stepBuilderFactory.get("workerStep")
        .<User, UserDto>chunk(100)
        .reader(databaseItemReader(null)) // 参数从 ExecutionContext 获取
        .processor(itemProcessor())
        .writer(itemWriter())
        .build();
}

// 3. 定义 Master Step
@Bean
public Step partitionedStep() {
    TaskExecutorPartitionHandler handler = new TaskExecutorPartitionHandler();
    handler.setStep(workerStep());
    handler.setTaskExecutor(taskExecutor());
    handler.setGridSize(4); // 4 个分片

    return stepBuilderFactory.get("masterStep")
        .partitioner("workerStep", partitioner())
        .handler(handler)
        .build();
}

8.4 使用线程池提升处理效率(TaskExecutor 配置)

配置项说明推荐值/策略
TaskExecutor 类型SimpleAsyncTaskExecutor(不复用线程,仅测试)、ThreadPoolTaskExecutor(标准线程池,生产首选)、ForkJoinPoolTaskExecutor(分治任务)。生产环境必须使用 ThreadPoolTaskExecutor
corePoolSize核心线程数。设置为 CPU 核心数(如 4, 8)。
maxPoolSize最大线程数。可设为 corePoolSize 的 2-4 倍(如 16-32)。
queueCapacity等待队列容量。有界队列推荐(如 100, 1000),防止资源耗尽。
rejectedExecutionHandler队列满时的拒绝策略。AbortPolicy(抛异常,默认)、CallerRunsPolicy(可降级防雪崩)、DiscardPolicy(静默丢弃)。
keepAliveSeconds非核心线程空闲存活时间。60 秒。
应用位置并行 Steps(split())、分片 Step(PartitionHandler)、异步 ItemProcessor/Writer。线程池是并行化的基础设施。
监控监控线程池的活跃线程数、队列大小、拒绝任务数。可通过 JMX 或 Micrometer 暴露指标。

8.5 大数据量下的分页读取与写入优化

优化点说明注意事项
ItemReader 分页读取避免 SELECT * 一次性加载所有数据导致 OOM。✅ 关键:排序键(sortKey)必须唯一且稳定,确保分页不漏/重。
ItemWriter 批量写入避免逐条 INSERT 性能低下。批量大小(batch_size)应与 chunkSize 匹配或为其倍数。
读写 chunkSize 匹配通常将 ItemReader 的 pageSize 设置为等于或略大于 chunkSize。减少数据库往返次数。
数据库优化写入表:临时删除非必要索引、外键约束,写完后重建;读取表:确保分页查询的 sortKey 有索引;连接池:配置足够大。写入前禁用索引是最有效的优化手段之一。
内存与 GC 优化避免在 ItemProcessor 中持有大数据引用;使用不可变对象;合理 JVM 参数(G1GC)。监控 GC 日志,避免频繁 Full GC。
流式处理 (Streaming)使用 HibernateCursorItemReader 或数据库原生流式 API,实现 O(1) 内存消耗。读取期间数据库连接必须保持,事务不能过早提交。

JdbcPagingItemReader 分页配置示例:

reader.setPageSize(1000); // 与 chunkSize 可不同
reader.setQueryProvider(new SqlPagingQueryProviderFactoryBean() {{
    setSelectClause("id, name, email");
    setFromClause("from users");
    setSortKey("id"); // 必须有唯一排序键
}}.getObject());

JPA/Hibernate 批量写入优化配置:

spring:
  jpa:
    properties:
      hibernate:
        jdbc:
          batch_size: 50
        order_inserts: true
        order_updates: true

性能优化总结:

  • 先调优 chunkSize:这是最基础且影响最大的参数。
  • 善用并行化:Parallel Steps 处理独立任务,Partitioning 处理海量数据。
  • 配置健壮的 TaskExecutor:线程池是并行化的引擎。
  • 分页读 + 批量写:大数据量的黄金组合。
  • 监控与测试:性能优化必须基于实际数据和监控指标,避免过度设计。

第9章 监控、日志与运维

9.1 使用 JobExplorer 和 JobOperator 查询执行状态

组件说明注意事项
JobExplorer只读访问 JobRepository 中的执行数据,用于查询和分析。线程安全,可安全用于监控服务;不会修改数据库状态。
JobOperator提供可操作的运维功能,可控制 Job 生命周期。改变系统状态,需谨慎使用(如 abort);方法可能抛 JobExecutionNotRunningException 等。

JobExplorer 核心方法:

  • getJobNames():获取所有 Job 名称
  • getJobInstances(String jobName, int start, int count):分页获取 Job 实例
  • getLastJobExecution(String jobName, JobParameters):获取最后一次执行
  • getJobExecutions(JobInstance):获取某实例的所有执行
  • getStepExecution(JobExecution, String stepName):获取 Step 执行详情

JobOperator 核心方法:

  • start(String jobName, String parameters):启动 Job
  • startNextInstance(String jobName):启动新实例(自动处理参数)
  • getRunningExecutions(String jobName):获取运行中的执行
  • abort(long executionId):终止执行(触发 STOPPED 状态)
  • stop(long executionId):请求停止(优雅终止)
  • getSummary(long executionId):获取执行摘要

JobExplorer 使用示例:

@Autowired
private JobExplorer jobExplorer;

public void analyzeJob(String jobName) {
    List<JobInstance> instances = jobExplorer.getJobInstances(jobName, 0, 10);
    for (JobInstance instance : instances) {
        List<JobExecution> executions = jobExplorer.getJobExecutions(instance);
        for (JobExecution execution : executions) {
            System.out.println("Status: " + execution.getStatus());
            System.out.println("StartTime: " + execution.getStartTime());
            System.out.println("ExitCode: " + execution.getExitStatus().getExitCode());
        }
    }
}

JobOperator 使用示例:

@Autowired
private JobOperator jobOperator;

public void stopJob(Long executionId) throws Exception {
    if (jobOperator.getRunningExecutions("dataExportJob").contains(executionId)) {
        jobOperator.stop(executionId); // 发送停止信号
    }
}

9.2 批处理执行数据的持久化与监控

数据类型存储位置监控价值
JOB_INSTANCEBATCH_JOB_INSTANCE统计某 Job 的总运行次数;检查是否已存在特定参数的实例。
JOB_EXECUTIONBATCH_JOB_EXECUTION查询最近 10 次执行状态;找出所有 FAILED 的执行进行分析。
STEP_EXECUTIONBATCH_STEP_EXECUTION分析哪个 Step 耗时最长;监控 readCount, writeCount, skipCount 是否正常。
EXECUTION_CONTEXTBATCH_JOB_EXECUTION_CTXBATCH_STEP_EXECUTION_CTX实现断点续传的精确位置;传递跨 Step 的业务参数。

监控实践:

  1. 定时报表:每日生成批处理成功率、平均耗时、数据量报表
  2. 告警:对 FAILED 状态、skipCount 异常升高、执行时间超阈值发送告警(邮件/钉钉/企业微信)
  3. 可视化:使用 Grafana + Prometheus 或 Kibana 展示执行趋势
  4. 审计:记录谁在何时启动了哪个 Job(可通过 JobParameters 记录 operator=user1)

建议为 JobRepository 表建立索引(如 JOB_EXECUTION.JOB_INSTANCE_ID, STEP_EXECUTION.JOB_EXECUTION_ID)。

9.3 集成 Actuator 实现健康检查与指标暴露

Actuator 端点说明监控指标示例
/actuator/health健康检查。默认检查 DataSource。status: UP/DOWN;若 JobRepository 连接失败,则 status: DOWN。
/actuator/metrics暴露 Micrometer 指标。spring.batch.job.execution.active:活跃 Job 执行数;spring.batch.job.execution.duration:Job 执行时长(直方图);spring.batch.step.execution.items:Step 处理项目数(计数器)。
/actuator/batch提供 Job/Step 的列表和详情。GET /actuator/batch/jobs 列出所有 Job;GET /actuator/batch/jobs/{name}/executions 列出某 Job 的执行。
自定义指标在 ItemProcessor 或 ItemWriter 中手动记录。将指标与业务逻辑结合,实现精细化监控。

Actuator 配置(application.yml):

management:
  endpoint:
    health:
      show-details: always
  endpoints:
    web:
      exposure:
        include: health,info,metrics,batch

Prometheus 集成配置:

management:
  endpoints:
    prometheus:
      enabled: true
  metrics:
    export:
      prometheus:
        enabled: true

自定义指标示例:

@Autowired
private MeterRegistry meterRegistry;

public T process(T item) {
    Timer.Sample sample = Timer.start(meterRegistry);
    try {
        // 处理逻辑
        T result = ...;
        sample.stop(Timer.builder("item.process.time").register(meterRegistry));
        return result;
    } catch (Exception e) {
        meterRegistry.counter("item.process.error",
            "type", e.getClass().getSimpleName()).increment();
        throw e;
    }
}

9.4 自定义监听器记录执行详情

监听器类型注解触发时机典型用途
Job 监听器@BeforeJob / @AfterJobJob 开始前 / 结束后。初始化资源、发送开始/结束通知、记录整体耗时和结果。
Step 监听器@BeforeStep / @AfterStepStep 开始前 / 结束后。记录 Step 级别耗时、检查 ExecutionContext、做 Step 级别的资源管理。
Chunk 监听器@BeforeChunk / @AfterChunk / @AfterChunkError每个 chunk 开始前 / 成功后 / 失败后。记录每个 chunk 的大小和耗时、实现细粒度重试或补偿逻辑、监控吞吐量波动。

Job 监听器示例:

@Component
public class JobLoggerListener {
    private static final Logger log = LoggerFactory.getLogger(JobLoggerListener.class);

    @BeforeJob
    public void beforeJob(JobExecution jobExecution) {
        log.info("Job [{}] started. Params: {}",
            jobExecution.getJobInstance().getJobName(),
            jobExecution.getJobParameters());
    }

    @AfterJob
    public void afterJob(JobExecution jobExecution) {
        log.info("Job [{}] finished with status: {}, exitCode: {}",
            jobExecution.getJobInstance().getJobName(),
            jobExecution.getStatus(),
            jobExecution.getExitStatus().getExitCode());

        if (jobExecution.getStatus() == BatchStatus.FAILED) {
            alertService.send("Job Failed: " +
                jobExecution.getJobInstance().getJobName());
        }
    }
}

Step 监听器示例:

@BeforeStep
public void beforeStep(StepExecution stepExecution) {
    log.info("Step [{}] starting. ReadCount: {}",
        stepExecution.getStepName(), stepExecution.getReadCount());
}

@AfterStep
public ExitStatus afterStep(StepExecution stepExecution) {
    log.info("Step [{}] ended. Status: {}, WriteCount: {}, SkipCount: {}",
        stepExecution.getStepName(),
        stepExecution.getStatus(),
        stepExecution.getWriteCount(),
        stepExecution.getSkipCount());
    return stepExecution.getExitStatus(); // 可修改 ExitStatus
}

Chunk 监听器示例:

@AfterChunk
public void afterChunk(ChunkContext context) {
    StepExecution stepExecution = context.getStepContext().getStepExecution();
    Chunk chunk = context.getChunk();
    log.debug("Chunk processed. Size: {}, WriteCount: {}",
        chunk.getItems().size(), stepExecution.getWriteCount());
}

@AfterChunkError
public void afterChunkError(ChunkContext context) {
    log.error("Chunk processing failed. Step: {}",
        context.getStepContext().getStepName());
    // 可在此处进行补偿或记录失败数据
}

监听器注册方式:

return stepBuilderFactory.get("step1")
    .<User, UserDto>chunk(10)
    .reader(reader())
    .processor(processor())
    .writer(writer())
    .listener(jobLoggerListener) // 注入监听器 Bean
    .listener(new CustomStepListener()) // 或直接 new
    .build();

9.5 日志记录最佳实践

实践原则说明推荐做法反模式
结构化日志使用 JSON 格式日志,便于机器解析和集中收集(如 ELK)。配合 logstash-logback-encoder 输出 JSON。字符串拼接:log.info("Job " + name + " started")
分层日志级别合理使用 DEBUG, INFO, WARN, ERROR。INFO: Job/Step 开始/结束;WARN: 可跳过的异常;ERROR: 致命错误;DEBUG: 详细流程(生产关闭)。所有日志都用 INFO,或在 ERROR 中打印堆栈但不抛出异常。
包含上下文信息每条日志应包含足够的上下文,便于追踪。在监听器中记录 JobExecutionId, StepExecutionId, JobParameters。只记录 “Processing started”,无法关联到具体执行。
避免记录敏感数据防止密码、身份证号等泄露。对敏感字段脱敏(如 email: user***@xxx.com)。log.info("Processing user: {}", user) 可能打印完整对象。
性能考量日志 I/O 可能成为瓶颈。使用异步日志(如 Logback AsyncAppender);DEBUG 级别在生产环境设为 OFF。ItemProcessor.process() 中每条记录都打 INFO 日志。
统一日志格式团队内约定日志模板。定义常量或工具类统一格式。每个人按自己习惯写日志,格式混乱。
集中化管理将日志发送到中心化系统。使用 Filebeat → Kafka → Logstash → Elasticsearch → Kibana 架构。日志分散在各服务器文件中。

结构化日志示例:

log.info("Job started",
    "jobName", jobExecution.getJobInstance().getJobName(),
    "jobId", jobExecution.getJobInstance().getInstanceId(),
    "params", jobExecution.getJobParameters().getParameters());

运维总结:

  • JobExplorer/JobOperator 是运维的”数据之眼”和”控制之手”。
  • Actuator + Micrometer + Prometheus + Grafana 是现代 Spring Boot 应用的标准监控栈。
  • 自定义监听器是实现业务定制化监控和告警的核心。
  • 结构化、有上下文、分级别的日志是故障排查的生命线。
  • 监控和日志的目标是:可观测性(Observability),即快速发现问题、定位根因、评估影响。

第10章 高级特性与扩展

10.1 条件流程控制(Flow、Decision、Transition)

概念说明注意事项
Flow将多个 Step 组织成一个可复用的执行单元,可作为 Job 的一部分或并行分支。Flow 可以嵌套,并可用于 split 实现并行。
JobExecutionDecider实现自定义决策逻辑,根据 JobExecution 状态决定下一步。决策结果(FlowExecutionStatus)用于匹配 Transition。
Transition(流转)定义执行路径的跳转规则,基于 ExitStatus。* 通配符匹配前缀;.end() 终止流程;.fail() 标记 Job 失败。
复合条件结合多个决策和流转实现复杂逻辑。流程定义清晰,避免过度复杂。

Flow 示例:

@Bean
public Flow dataProcessingFlow() {
    return new FlowBuilder<SimpleFlow>("dataProcessingFlow")
        .start(step1())
        .next(step2())
        .build();
}

@Bean
public Job conditionalJob() {
    return jobBuilderFactory.get("conditionalJob")
        .start(dataProcessingFlow())
        .next(decisionStep())
        .on("COMPLETED").to(finalStep())
        .on("FAILED").stop()
        .build();
}

JobExecutionDecider 示例:

@Component
public class FileExistsDecider implements JobExecutionDecider {
    @Override
    public FlowExecutionStatus decide(JobExecution jobExecution,
                                       StepExecution stepExecution) {
        String fileName = jobExecution.getJobParameters().getString("input.file");
        boolean exists = Files.exists(Paths.get(fileName));
        return exists ? FlowExecutionStatus.COMPLETED : FlowExecutionStatus.FAILED;
    }
}

// 在 Job 中使用:
.next(fileExistsDecider())
    .on("COMPLETED").to(importStep())
    .on("FAILED").to(skipImportStep())

Transition 流转示例:

return jobBuilderFactory.get("retryJob")
    .start(importStep())
    .on("FAILED*").to(errorHandlingStep()) // 匹配 FAILED 开头的状态
    .from(importStep())
    .on("COMPLETED_WITH_SKIPS").to(qualityReviewStep())
    .end()
    .build();

10.2 Job 的嵌套与子流程调用

方式说明注意事项
JobStep将一个 Job 作为另一个 Job 的一个 Step 执行,实现 Job 嵌套。使用 JobStepBuilder 创建 JobStep,在主 Job 中引用。
Flow 复用将公共步骤序列定义为 Flow,在多个 Job 中引用。更轻量级的复用方式,推荐优先使用。

JobStep 示例:

// 子 Job
@Bean
public Job subJob() { ... }

// JobStep
@Bean
public Step subJobStep(Job subJob) {
    return jobStepBuilderFactory.get("subJobStep")
        .job(subJob)
        .parametersExtractor((jobExecution, stepExecution) ->
            new JobParametersBuilder(jobExecution.getJobParameters())
                .addString("source", "parentJob") // 可传递/转换参数
                .toJobParameters())
        .build();
}

// 主 Job
@Bean
public Job masterJob() {
    return jobBuilderFactory.get("masterJob")
        .start(preparationStep())
        .next(subJobStep(subJob())) // 调用子 Job
        .next(finalizationStep())
        .build();
}

10.3 使用 Spring Expression Language(SpEL)控制流程

用途SpEL 示例注意事项
动态参数JobParametersBuilder 中使用 SpEL 生成参数值。需在运行时解析。
条件流转.filter("T(java.lang.Math).random() < 0.5 ? 'A' : 'B'")实现 A/B 测试或随机分流。
@Value 注解注入@Value("#{jobParameters['input.file']}")Bean 需在 Job 启动时创建(@StepScope@JobScope)。

⚠️ 作用域 (Scope) 关键:

  • @StepScope:Bean 在 Step 执行时创建,可访问 StepExecutionContext 和 JobParameters。
  • @JobScope:Bean 在 Job 执行时创建。
@Component
@StepScope // 必须声明作用域
public class ScopedItemReader implements ItemReader<Foo> {
    @Value("#{jobParameters['fileName']}")
    private String fileName;
    // ...
}

10.4 远程分片(Remote Chunking)简介

概念说明优缺点
远程分片 (Remote Chunking)Master 负责读取数据并通过消息中间件发送给多个 Worker 进行处理和写入,实现计算与 I/O 的物理分离。✅ 优点:Master 可专注于高效读取;Worker 可水平扩展;解耦读取和处理。❌ 缺点:架构复杂;网络延迟和序列化开销。

架构与流程:

  1. Master:ItemReader 读取数据 → MessagingTemplate 将 chunk 发送到 Message Channel
  2. Message Middleware:传递数据块(Kafka, RabbitMQ, JMS)
  3. Worker(s):监听 Message Channel → 接收 chunk → 执行 ItemProcessor/ItemWriter → 回传结果

技术栈:

  • Master: Spring Integration + MessagingTemplate
  • Middleware: Kafka, RabbitMQ, JMS
  • Worker: Message-Driven POJO (MDP) 或 @StreamListener
  • 序列化: JSON, Avro, Protobuf

10.5 与消息队列集成(如 Kafka、RabbitMQ)

集成方式说明技术方案典型场景
作为远程分片的通道见 10.4 节。Spring Integration + Spring Batch Remote Chunking。大规模数据处理,需要解耦读写。
作为 ItemReader 数据源从消息队列消费数据作为批处理输入。KafkaItemReader(自定义或社区库)、AmqpItemReader处理积压的消息(如补偿任务、报表生成)。
作为 ItemWriter 输出目标将处理结果发送到消息队列。KafkaItemWriterAmqpItemWriter批处理结果触发后续微服务处理。
事件通知Job/Step 状态变更时发送消息。在监听器中使用 RabbitTemplate/KafkaTemplate 发送状态事件。解耦监控系统,实现事件驱动架构。
使用 Spring Cloud Stream统一抽象消息中间件。@StreamListener + @ServiceActivator实现技术无关的消息集成,推荐使用。

10.6 与调度框架集成(如 Quartz、Scheduled Tasks)

方式说明优缺点
Spring @ScheduledSpring 内置的轻量级定时任务。✅ 优点:简单,无需额外依赖。❌ 缺点:不支持集群,功能简单。
Quartz功能强大的企业级调度框架,支持集群和复杂调度。✅ 优点:高可用、复杂调度、集群环境。
Kubernetes CronJob在云原生环境中,使用 K8s 的 CronJob 调度批处理 Pod。✅ 优点:云原生标准,天然支持水平扩展和资源隔离。❌ 缺点:依赖 K8s 环境。

Spring @Scheduled 示例:

@Component
public class JobScheduler {
    @Autowired
    private JobLauncher jobLauncher;
    @Autowired
    private Job reportJob;

    @Scheduled(cron = "0 0 2 * * ?") // 每天凌晨2点
    public void runDailyReport() {
        try {
            JobParameters params = new JobParametersBuilder()
                .addDate("runDate", new Date())
                .toJobParameters();
            jobLauncher.run(reportJob, params);
        } catch (Exception e) {
            log.error("Failed to run job", e);
        }
    }
}

Quartz 集成示例:

@Component
public class QuartzJobLauncher implements Job {
    @Override
    public void execute(JobExecutionContext context) throws JobExecutionException {
        JobLauncher jobLauncher = (JobLauncher) context.getScheduler()
            .getContext().get("jobLauncher");
        Job job = (Job) context.getJobDetail().getJobDataMap().get("job");

        try {
            JobParameters params = new JobParametersBuilder()
                .addLong("fireTime", context.getFireTime().getTime())
                .toJobParameters();
            JobExecution execution = jobLauncher.run(job, params);
        } catch (Exception e) {
            throw new JobExecutionException(e);
        }
    }
}

// 配置
@Bean
public JobDetail quartzJobDetail() {
    return JobBuilder.newJob(QuartzJobLauncher.class)
        .usingJobData("job", reportJob)
        .storeDurably()
        .build();
}

@Bean
public Trigger quartzTrigger() {
    return TriggerBuilder.newTrigger()
        .forJob(quartzJobDetail())
        .withSchedule(CronScheduleBuilder.cronSchedule("0 0 2 * * ?"))
        .build();
}

Kubernetes CronJob 示例:

apiVersion: batch/v1
kind: CronJob
metadata:
  name: daily-report
spec:
  schedule: "0 2 * * *"
  jobTemplate:
    spec:
      template:
        spec:
          containers:
            - name: batch-app
              image: my-batch-app:latest
              args:
                - "--spring.batch.job.names=dailyReportJob"
          restartPolicy: OnFailure

调度选择建议:

  • 单机简单任务:@Scheduled
  • 高可用、复杂调度、集群环境:Quartz
  • 云原生环境:Kubernetes CronJob