第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 是可重启的单位。 |
| Step | Job 中的一个独立阶段,通常对应一个读-处理-写操作序列。 | 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 等。 | 用于判断执行结果,决定是否重试或继续。 |
| ExitStatus | Step 的退出状态,可自定义(如 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。 |
| Gradle | 在 build.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 | 控制元数据表初始化行为:ALWAYS、EMBEDDED、NEVER。 | ALWAYS:每次启动都建表(生产慎用);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 表示每批处理的记录数。 |
| Job | Job job = jobBuilderFactory.get("jobName").start(step).build(); | 定义一个包含 Step 的 Job。 | 见上例。 | Job 名称必须唯一。 |
| Step | Step step = stepBuilderFactory.get("stepName")...build(); | 定义一个处理步骤。 | 见上例。 | 一个 Job 至少包含一个 Step。 |
2.5 运行并观察批处理执行日志
| 日志内容 | 说明 | 注意事项 |
|---|
| Application Started | Spring Boot 应用启动完成,批处理自动执行(若未禁用)。 | 默认情况下,Spring Batch 在应用启动时自动运行已定义的 Job。 |
| Starting Job | 日志显示 Job 开始执行,包含 Job 名称和 JobParameters。 | 可通过日志确认 Job 是否按预期启动。 |
| Reading from ItemReader | 显示读取数据的日志(取决于日志级别)。 | 可通过调试模式查看每条记录的读取情况。 |
| Processed Item | 显示 ItemProcessor 处理结果(如有日志输出)。 | 可用于验证数据转换逻辑是否正确。 |
| Writing Items | 显示写入数据的数量和内容。 | 确认写入是否成功,chunk 提交是否正常。 |
| Finished Job | Job 执行完成,输出最终状态(COMPLETED、FAILED 等)。 | 成功状态应为 COMPLETED,失败则需排查异常堆栈。 |
| 异常堆栈(Exception Stack Trace) | 若执行失败,会输出详细异常信息。 | 重点关注 Caused by 部分,定位根本原因。 |
| Batch Metadata Logs | 记录 JobExecution、StepExecution 的持久化操作。 | 表明元数据已写入数据库,支持后续重启。 |
第3章 Job 与 Step 的定义与控制
3.1 Job 的创建与配置(JobBuilderFactory)
| 方法名 | 语法 | 用途说明 | 代码示例 | 注意事项 |
|---|
get | JobBuilder get(String jobName) | 获取一个 JobBuilder 实例,用于构建 Job。 | Job job = jobBuilderFactory.get("importUserJob").start(step1).build(); | jobName 必须唯一,否则会覆盖已有 Job。 |
start | JobBuilder start(Step firstStep) | 指定 Job 的第一个 Step。 | jobBuilderFactory.get("job").start(step1).build(); | 每个 Job 至少要有一个 Step。 |
next | JobBuilder next(Step nextStep) | 添加后续 Step,实现线性流程。 | jobBuilderFactory.get("job").start(step1).next(step2).build(); | Step 执行顺序由调用链决定。 |
flow | JobBuilder flow(Step step) | 将 Step 包装为 Flow,用于复杂流程控制。 | jobBuilderFactory.get("job").start(flow1).build(); | 需结合 end()、on().to() 等方法实现条件跳转。 |
split | JobBuilder split(TaskExecutor taskExecutor) | 创建并行分支(Parallel Steps)。 | jobBuilderFactory.get("job").split(taskExecutor).add(flow1, flow2); | 需提供 TaskExecutor 实现并发执行。 |
incrementer | JobBuilder incrementer(JobParametersIncrementer) | 设置参数增量器,避免 JobInstance 冲突。 | jobBuilderFactory.get("job").incrementer(new DailyJobIncrementer()).start(step); | 常用于每日执行的 Job,自动添加时间戳参数。 |
listener | JobBuilder listener(JobExecutionListener listener) | 注册 Job 级别的监听器。 | jobBuilderFactory.get("job").listener(jobListener).start(step); | 可多次调用添加多个监听器。 |
preventRestart | JobBuilder preventRestart(boolean preventRestart) | 控制 Job 是否可重启。 | jobBuilderFactory.get("job").preventRestart(true).start(step); | 默认 false(可重启),设为 true 则失败后不能重启。 |
repository | JobBuilder repository(JobRepository jobRepository) | 指定使用的 JobRepository。 | 一般无需手动设置,由 Spring 自动注入。 | 高级用法,用于自定义元数据存储。 |
build | Job build() | 构建并返回最终的 Job 实例。 | Job job = jobBuilderFactory.get("job").start(step).build(); | 必须调用 build() 才能生成 Job 对象。 |
3.2 Step 的创建与配置(StepBuilderFactory)
| 方法名 | 语法 | 用途说明 | 代码示例 | 注意事项 |
|---|
get | StepBuilder get(String stepName) | 获取 StepBuilder 实例。 | Step step = stepBuilderFactory.get("processDataStep").<String, String>chunk(10).reader(reader).writer(writer).build(); | stepName 在 Job 内应唯一。 |
chunk | SimpleStepBuilder chunk(int chunkSize) | 配置基于块(chunk)的处理模式。 | stepBuilderFactory.get("step").<Integer, Integer>chunk(50).reader(itemReader).processor(itemProcessor).writer(itemWriter); | chunkSize 表示每批处理多少条记录后提交事务。 |
tasklet | TaskletStepBuilder tasklet(Tasklet tasklet) | 配置 Tasklet 模式的 Step(非 chunk)。 | stepBuilderFactory.get("cleanUpStep").tasklet(cleanupTasklet).build(); | 适用于无需读-处理-写循环的任务,如清理、通知等。 |
reader | SimpleStepBuilder reader(ItemReader reader) | 设置数据读取器。 | 见 chunk 示例。 | 必须与 chunk 模式配合使用。 |
processor | SimpleStepBuilder processor(ItemProcessor processor) | 设置数据处理器。 | 见 chunk 示例。 | 可选,若无需处理可省略。 |
writer | SimpleStepBuilder writer(ItemWriter writer) | 设置数据写入器。 | 见 chunk 示例。 | 必须与 chunk 模式配合使用。 |
listener | StepBuilder listener(StepExecutionListener listener) | 注册 Step 级别的监听器。 | stepBuilderFactory.get("step").listener(stepListener).chunk(10)... | 可多次调用添加多个监听器。 |
transactionManager | StepBuilder transactionManager(PlatformTransactionManager tm) | 指定事务管理器。 | 通常使用默认事务管理器,除非有特殊需求。 | 在分布式事务场景下可能需要自定义。 |
allowStartIfComplete | StepBuilder allowStartIfComplete(boolean allow) | 是否允许 Step 已完成时仍可启动。 | stepBuilderFactory.get("step").allowStartIfComplete(true).chunk(10)... | 默认 false,设为 true 可重复执行已完成的 Step。 |
startLimit | StepBuilder startLimit(int startLimit) | 设置 Step 最大启动次数。 | stepBuilderFactory.get("step").startLimit(3).chunk(10)... | 超过限制将抛出 StartLimitExceededException。 |
build | Step 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 添加增量值(如时间戳),避免参数重复。 | 常见实现:RunIdIncrementer、DailyJobTimeIncrementer。 |
3.6 Job 的监听器(JobExecutionListener)
| 方法名 | 语法 | 用途说明 | 代码示例 | 注意事项 |
|---|
beforeJob | void beforeJob(JobExecution jobExecution) | Job 开始前执行的逻辑。 | 见下方代码块 | 可用于初始化资源、记录日志、发送通知等。 |
afterJob | void 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)
| 方法名 | 语法 | 用途说明 | 代码示例 | 注意事项 |
|---|
beforeStep | void beforeStep(StepExecution stepExecution) | Step 开始前执行。 | 见下方代码块 | 可用于初始化 Step 相关资源。 |
afterStep | ExitStatus 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 或数据库特定分页语法),适合大数据集。 | 避免长事务,支持并行处理 |
StaxEventItemReader | XML 文件 | 基于 SAX 的事件驱动解析,逐个读取 XML 元素,内存友好。 | 结构化 XML 数据导入 |
RepositoryItemReader | JPA Repository | 调用 Spring Data JPA 的 PagingAndSortingRepository 分页读取实体。 | 已有 JPA 实体和 Repository 的项目 |
ListItemReader | Java 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 场景 |
JpaItemWriter | JPA 实体管理 | 通过 EntityManager 批量持久化实体,支持级联。 | 已使用 JPA 的项目,需复杂对象图保存 |
CompositeItemWriter | 多个目标 | 包装多个 ItemWriter,将同一数据写入多个目的地。 | 数据同步、备份、多格式输出 |
ClassifierCompositeItemWriter | 分类写入 | 根据分类器(Classifier)将不同数据路由到不同 ItemWriter。 | 条件分发,如按地区写入不同表/文件 |
ItemWriterAdapter | 任意服务方法 | 将普通服务方法适配为 ItemWriter,通过反射调用。 | 调用遗留系统、外部 API 写入数据 |
MongoItemWriter | MongoDB | 写入文档到 MongoDB 集合。 | NoSQL 场景,JSON 数据存储 |
KafkaItemWriter | Kafka 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_size 和 hibernate.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_INSTANCE | BATCH_JOB_INSTANCE 表 | 统计某 Job 的总运行次数;检查是否已存在特定参数的实例。 |
| JOB_EXECUTION | BATCH_JOB_EXECUTION 表 | 查询最近 10 次执行状态;找出所有 FAILED 的执行进行分析。 |
| STEP_EXECUTION | BATCH_STEP_EXECUTION 表 | 分析哪个 Step 耗时最长;监控 readCount, writeCount, skipCount 是否正常。 |
| EXECUTION_CONTEXT | BATCH_JOB_EXECUTION_CTX、BATCH_STEP_EXECUTION_CTX | 实现断点续传的精确位置;传递跨 Step 的业务参数。 |
监控实践:
- 定时报表:每日生成批处理成功率、平均耗时、数据量报表
- 告警:对 FAILED 状态、skipCount 异常升高、执行时间超阈值发送告警(邮件/钉钉/企业微信)
- 可视化:使用 Grafana + Prometheus 或 Kibana 展示执行趋势
- 审计:记录谁在何时启动了哪个 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 / @AfterJob | Job 开始前 / 结束后。 | 初始化资源、发送开始/结束通知、记录整体耗时和结果。 |
| Step 监听器 | @BeforeStep / @AfterStep | Step 开始前 / 结束后。 | 记录 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 可水平扩展;解耦读取和处理。❌ 缺点:架构复杂;网络延迟和序列化开销。 |
架构与流程:
- Master:ItemReader 读取数据 → MessagingTemplate 将 chunk 发送到 Message Channel
- Message Middleware:传递数据块(Kafka, RabbitMQ, JMS)
- 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 输出目标 | 将处理结果发送到消息队列。 | KafkaItemWriter、AmqpItemWriter。 | 批处理结果触发后续微服务处理。 |
| 事件通知 | Job/Step 状态变更时发送消息。 | 在监听器中使用 RabbitTemplate/KafkaTemplate 发送状态事件。 | 解耦监控系统,实现事件驱动架构。 |
| 使用 Spring Cloud Stream | 统一抽象消息中间件。 | @StreamListener + @ServiceActivator | 实现技术无关的消息集成,推荐使用。 |
10.6 与调度框架集成(如 Quartz、Scheduled Tasks)
| 方式 | 说明 | 优缺点 |
|---|
Spring @Scheduled | Spring 内置的轻量级定时任务。 | ✅ 优点:简单,无需额外依赖。❌ 缺点:不支持集群,功能简单。 |
| 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