Spring Batch批处理核心原理:Chunk机制、重启策略与资源隔离
1. 为什么Spring Batch不是“另一个定时任务框架”——从真实业务场景切入
我第一次在生产环境里碰上Spring Batch,是在一个电商订单对账系统里。当时团队用@Scheduled写了个每5分钟跑一次的定时任务,逻辑是“查出昨天所有未对账订单,逐条调用第三方支付接口核验状态”。上线两周后,数据库连接池频繁告警,日志里全是Connection timeout,运维同事半夜打电话问我:“你那个‘小脚本’是不是偷偷开了300个线程?”——其实它根本没开线程,只是单线程循环处理27万条订单,每条调用一次HTTP,平均耗时1.8秒,整个任务跑完要13个小时,而下一轮调度又开始了,线程全卡在等待HTTP响应上。
这就是典型误用:把批处理当成了“带循环的定时任务”。Spring Batch真正的价值,从来不是“让批量操作能跑起来”,而是让百万级数据在可控、可监控、可恢复、不拖垮系统的前提下,稳稳落地。它不解决“要不要跑”的问题,它解决的是“怎么跑才不死、跑错了怎么救、跑一半断电了怎么办、跑慢了怎么提速”这些真正折磨人的细节。
关键词里反复出现的“示例”,恰恰暴露了当前学习者的普遍困境:网上90%的Spring Batch教程,只教你怎么写一个能编译通过的Job,却从不告诉你——
- 当ItemReader读到第8万条记录时OOM了,堆栈里显示org.springframework.batch.item.support.builder.BuilderSupport$BuilderSupport,这到底是谁的锅?
- Step执行失败后,重启Job为什么从头开始而不是接着上次断点继续?
- 同一个Job配置里同时用了JdbcCursorItemReader和JdbcPagingItemReader,结果分页SQL被自动改写成错误语法,报错信息里连表名都看不清;
- 测试环境跑得好好的,一上生产就卡在TaskExecutor初始化阶段,日志里只有“Initializing ExecutorService”,再无下文。
这些不是边缘case,而是每天都在发生的现实。所以这篇内容不叫“SpringBoot整合Spring Batch入门”,它叫**《SpringBoot整合Spring Batch:从能跑通到敢上线的完整链路》**。我会带着你,从零搭建一个真实可用的批处理模块,重点拆解那些官方文档里一笔带过、但线上踩坑后必须花三天才能定位的硬核细节。核心关键词就三个:Chunk机制、重启策略、资源隔离——它们才是Spring Batch区别于其他框架的生死线。
2. Chunk不是“分块”,而是批处理的原子性契约——原理与陷阱
很多人看到“Chunk-oriented processing”就理解为“把大数据切成小块处理”,这没错,但远远不够。Chunk的本质,是一份事务边界+重试契约+状态快照协议。它规定:在一个Chunk内,所有read→process→write操作必须被视为一个不可分割的单元。要么全部成功提交,要么全部回滚,且失败时必须能精确还原到Chunk开始前的状态。
2.1 Chunk生命周期的四个不可跳过阶段
我们以一个典型订单对账Job为例,其Step配置如下:
@Bean public Step orderReconciliationStep( ItemReader<Order> reader, ItemProcessor<Order, ReconciliationResult> processor, ItemWriter<ReconciliationResult> writer, PlatformTransactionManager transactionManager) { return stepBuilderFactory.get("orderReconciliationStep") .<Order, ReconciliationResult>chunk(1000) // 关键:Chunk size=1000 .reader(reader) .processor(processor) .writer(writer) .faultTolerant() // 启用容错 .skipPolicy(new AlwaysSkipItemSkipPolicy()) // 跳过异常项 .retryLimit(3) // 重试3次 .retry(Exception.class) // 对所有Exception重试 .transactionManager(transactionManager) .build(); }这个chunk(1000)背后,Spring Batch实际执行的是四阶段流水线:
Read Phase(读取阶段):ItemReader连续调用read()方法,直到获取1000条Order对象,或返回null(表示数据源结束)。注意:此时不触发任何数据库事务,读取过程完全独立于后续写入。
Process Phase(处理阶段):对读取到的1000条Order,逐条调用ItemProcessor.process()。若某条Order处理抛出异常(如空指针),且该异常被retry或skip策略捕获,则此条记录被标记为“跳过”或“重试”,但Chunk继续执行,直到凑满1000条有效处理结果(跳过/重试后的剩余条目)。
Write Phase(写入阶段):将1000条处理后的ReconciliationResult对象,一次性交给ItemWriter.write()。此时事务正式开启,write()方法内部必须完成所有数据库INSERT/UPDATE操作,并在方法返回前提交事务。如果write()中途抛出异常(如唯一键冲突),整个Chunk回滚,已处理的999条记录全部撤销。
Commit Phase(提交阶段):事务成功提交后,Spring Batch立即向JobRepository(通常是数据库表BATCH_STEP_EXECUTION_CONTEXT)写入本次Chunk的执行上下文,包括:
READ_COUNT=1000(本次读取条数)WRITE_COUNT=1000(本次写入条数)COMMIT_COUNT=1(本次提交次数)STEP_NAME="orderReconciliationStep"KEY="orderReconciliationStep:1"(唯一标识)
提示:这个上下文表就是重启能力的基石。Job重启时,Batch会查询此表,找到最后一次成功commit的KEY,然后让ItemReader从对应位置继续读取(如JdbcCursorItemReader会重置游标到第1001条)。
2.2 Chunk Size选型:不是越大越好,也不是越小越稳
网上教程常建议“设为100或1000”,但真实场景中,这个值必须根据三要素动态计算:
| 要素 | 影响逻辑 | 实测案例 |
|---|---|---|
| 单条数据处理内存占用 | Chunk内所有对象驻留在JVM堆中,直到write()完成。若单条Order对象含10个String字段+3个BigDecimal+1个List ,实测占内存约12KB,则1000条≈12MB | 某金融对账系统,Chunk=500时Full GC频次为2min/次;调至200后降为15min/次 |
| 数据库事务锁粒度 | MySQL InnoDB对INSERT语句加行锁,但大批量INSERT可能升级为间隙锁(Gap Lock)。Chunk=5000时,对账表锁等待超时率达37%;降至1000后降至1.2% | 生产环境监控发现,锁等待时间与Chunk size呈近似平方关系 |
| 网络IO稳定性 | write()调用外部API(如支付核验)时,Chunk过大导致单次HTTP请求超时风险陡增。某第三方接口SLA为99.5%,Chunk=1000时单次失败概率≈0.5%,Chunk=100时≈0.05% | 通过Prometheus监控write阶段失败率,发现拐点在Chunk=300 |
我的经验公式:推荐ChunkSize = min(200, floor(可用堆内存×0.3 ÷ 单条对象内存))
其中0.3是保守系数(预留GC空间),单条对象内存可通过JProfiler采样获得。例如:
- JVM堆设为2GB → 可用堆≈1.6GB(-Xmx2g -XX:MaxMetaspaceSize=256m)
- 单条Order对象内存≈8KB
- 计算:floor(1.6×1024×1024×1024 × 0.3 ÷ 8192) ≈ floor(62914560) ≈ 62914560 → 显然不合理
- 实际应结合IO瓶颈:若write()调用外部API平均耗时200ms,目标吞吐量5000条/分钟,则单Chunk耗时需≤12s → 12s÷0.2s=60条 → 最终选定Chunk=50
注意:这个值必须在预发布环境用真实数据压测验证。我曾因直接套用测试环境参数,导致生产Job启动后CPU持续100%,排查发现是Chunk=1000时JDBC驱动缓存溢出,切换为HikariCP并启用
cachePrepStmts=true后解决。
2.3 容错机制的致命误区:Skip与Retry的混用灾难
faultTolerant()开启后,开发者常犯两个错误:
错误1:对同一异常同时配置skip和retry
.faultTolerant() .skip(Exception.class) // 跳过所有异常 .skipLimit(1000) .retry(Exception.class) // 又重试所有异常 .retryLimit(3)这会导致:当第1条记录抛出IOException,Batch先尝试重试3次,失败后计入skip计数;第2条记录同样IOException,同样重试3次再skip……最终skipLimit=1000很快耗尽,Job强制终止。正确做法是明确区分场景:
skip用于业务可容忍的脏数据(如订单金额为负数,直接跳过不处理)retry用于临时性故障(如网络抖动、数据库短暂不可用)
错误2:Retry策略未限定具体异常类型
.retry(RuntimeException.class) // 看似合理但RuntimeException包含NullPointerException、IllegalArgumentException等编程错误,重试毫无意义。应精确到:
.retry(ConnectException.class) // 网络连接异常 .retry(SQLTimeoutException.class) // 数据库超时 .retry(IOException.class) // IO异常 .retryLimit(3)更关键的是,retry必须配合退避策略(BackOff Policy),否则重试会雪崩:
.retry(ConnectException.class) .retryLimit(3) .backOffPolicy(new ExponentialBackOffPolicy() {{ setInitialInterval(1000L); // 首次重试间隔1秒 setMultiplier(2.0); // 每次翻倍 setMaxInterval(10000L); // 最大间隔10秒 }});3. Job重启不是“重新运行”,而是状态机驱动的精准续跑——从源码看执行流程
Spring Batch的Job重启能力,常被误解为“点一下Restart按钮就从头再来”。实际上,它是一个严格的状态机(State Machine),其核心在于StepExecution的ExecutionContext持久化与恢复。我们来看Job重启时的真实执行路径。
3.1 JobRepository:所有状态的唯一真相源
JobRepository是Spring Batch的“大脑”,默认实现为JdbcJobRepository,它依赖5张核心表:
BATCH_JOB_INSTANCE:记录Job每次执行的唯一实例(job_name + job_params hash)BATCH_JOB_EXECUTION:记录每次Job执行的总体状态(STARTED/COMPLETED/FAILED)BATCH_STEP_EXECUTION:记录每个Step的执行详情(READ_COUNT/WRITE_COUNT等)BATCH_JOB_EXECUTION_PARAMS:存储Job启动参数(如--input.file=path.csv)BATCH_STEP_EXECUTION_CONTEXT:最关键,存储Step执行过程中的临时状态(如游标位置、已处理ID列表)
当Job首次运行时,BATCH_STEP_EXECUTION_CONTEXT中会写入:
{ "stepName": "orderReconciliationStep", "read.count": 1000, "write.count": 1000, "commit.count": 1, "cursor.position": 1000, "last.processed.id": "ORD_20231001_000001" }3.2 Restart触发的三步校验链
当你调用JobOperator.restart(executionId)时,Batch执行以下校验:
Step 1:JobInstance合法性检查
查询BATCH_JOB_INSTANCE,确认该JobInstance存在且JOB_INSTANCE_ID匹配。若不存在,抛出NoSuchJobException。
Step 2:StepExecution状态校验
查询BATCH_STEP_EXECUTION,要求:
STATUS必须为FAILED、STOPPED或ABANDONED(不能是COMPLETED)VERSION字段必须为最新(防并发修改)- 若
STATUS=COMPLETED,则不允许restart,只能start()新实例
Step 3:ExecutionContext恢复与游标重置
这是最易出错的环节。以JdbcCursorItemReader为例,其open()方法会:
- 从
BATCH_STEP_EXECUTION_CONTEXT读取cursor.position值(如1000) - 执行原始SQL的
LIMIT子句,但不是简单加OFFSET,而是:
这种写法确保即使中间有记录被删除,也能准确定位到第1001条。SELECT * FROM orders WHERE id > (SELECT id FROM orders ORDER BY id LIMIT 1 OFFSET 1000) ORDER BY id LIMIT 1000
注意:若你自定义ItemReader未实现
update()方法(用于保存游标),则重启时永远从头开始。我曾遇到一个自研MongoDB Reader,因忘记在update()中写入lastProcessedId,导致每次重启都重复处理全部数据。
3.3 生产环境必须关闭的“安全开关”:allowStartIfComplete
Spring Batch默认禁止对COMPLETED状态的Job再次启动,防止重复执行。但某些场景(如每日对账Job需手动重跑某天数据)需要绕过此限制。配置方式:
@Bean public Job orderReconciliationJob() { return jobBuilderFactory.get("orderReconciliationJob") .start(orderReconciliationStep()) .listener(jobExecutionListener()) // 自定义监听器 .incrementer(new RunIdIncrementer()) // 每次生成新JobInstance .preventRestart() // 默认开启,禁止restart // 若要允许重跑,注释掉上一行,改为: // .allowStartIfComplete(true) // ⚠️ 生产慎用! .build(); }⚠️ 重大风险提示:allowStartIfComplete=true会使Job无视BATCH_JOB_INSTANCE的唯一性约束,导致同一组参数多次生成JobInstance。若你的Job逻辑未做幂等性设计(如未校验BATCH_JOB_EXECUTION_PARAMS中的日期参数),将引发数据重复写入。我的解决方案是:
- 在Job启动时,通过
JobParameters注入run.date=20231001 - 在ItemProcessor中,校验待处理订单的
create_date是否等于run.date,否则直接skip - 同时在JobListener的
beforeJob()中,查询BATCH_JOB_EXECUTION确认该日期是否已成功执行,若已存在则抛出JobExecutionAlreadyRunningException
4. 资源隔离:为什么你的Batch Job总拖垮主业务——线程池与事务管理实战
Spring Boot默认将所有Bean注册到同一个ApplicationContext,而Spring Batch的JobLauncher、TaskExecutor、TransactionManager若未显式隔离,极易与Web MVC的线程池、事务管理器产生资源争抢。我见过最典型的事故:一个每小时跑一次的库存同步Job,因共用Tomcat的tomcatThreadPool,导致用户下单接口TP99从120ms飙升至2.3s。
4.1 TaskExecutor:Batch专属线程池的硬性配置
Spring Batch的JobLauncher默认使用SimpleAsyncTaskExecutor,它不复用线程,每次new Thread(),在高并发Job场景下会创建海量线程,迅速耗尽系统资源。必须替换为ThreadPoolTaskExecutor:
@Configuration public class BatchConfig { @Bean @Primary // ⚠️ 关键:标记为主Bean,避免与其他TaskExecutor冲突 public TaskExecutor batchTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(4); // 核心线程数=CPU核数 executor.setMaxPoolSize(8); // 最大线程数=2×CPU核数 executor.setQueueCapacity(100); // 任务队列容量 executor.setThreadNamePrefix("batch-"); // 线程名前缀,便于日志追踪 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } @Bean public JobLauncher jobLauncher( @Qualifier("batchTaskExecutor") TaskExecutor taskExecutor, JobRepository jobRepository) { SimpleJobLauncher jobLauncher = new SimpleJobLauncher(); jobLauncher.setJobRepository(jobRepository); jobLauncher.setTaskExecutor(taskExecutor); // 绑定专属线程池 return jobLauncher; } }为什么CallerRunsPolicy是最佳选择?
当队列满时,CallerRunsPolicy会让调用线程(即Web请求线程)自己执行任务。这看似降低吞吐,实则起到熔断保护作用:若Batch任务积压,Web接口响应变慢,自然减少新请求涌入,避免系统雪崩。相比AbortPolicy(直接丢弃)或DiscardOldestPolicy(丢弃最老任务),它更符合生产环境的稳定性诉求。
4.2 TransactionManager:Batch与Web事务必须物理隔离
Spring Boot默认的DataSourceTransactionManager会被Web层和Batch层共享。问题在于:
- Web层事务通常短(毫秒级),Batch层事务长(分钟级)
- 若Batch Step的事务未提交,Web层的
@Transactional方法可能因连接池耗尽而超时
解决方案:为Batch创建独立DataSource
@Configuration public class DataSourceConfig { @Bean @ConfigurationProperties("spring.datasource.batch") // 读取application.yml中batch数据源 public DataSource batchDataSource() { return DataSourceBuilder.create().build(); } @Bean public PlatformTransactionManager batchTransactionManager( @Qualifier("batchDataSource") DataSource dataSource) { return new DataSourceTransactionManager(dataSource); } @Bean public JobRepository jobRepository( @Qualifier("batchDataSource") DataSource dataSource, @Qualifier("batchTransactionManager") PlatformTransactionManager transactionManager) throws Exception { JobRepositoryFactoryBean factory = new JobRepositoryFactoryBean(); factory.setDataSource(dataSource); factory.setTransactionManager(transactionManager); factory.setDatabaseType("mysql"); return factory.getObject(); } }对应的application.yml配置:
spring: datasource: # 主数据源(Web业务) url: jdbc:mysql://localhost:3306/main_db username: web_user password: xxx datasource: batch: url: jdbc:mysql://localhost:3306/batch_db # 独立库,或同库不同schema username: batch_user password: xxx提示:
batch_db不必是物理独立库,可以是同一MySQL实例下的不同schema,但必须保证batch_user权限仅限该schema,避免Batch Job意外修改主业务表。
4.3 内存泄漏的隐形杀手:ItemReader/Writer的静态缓存
很多开发者为提升性能,在ItemReader中使用静态Map缓存字典数据:
public class OrderReader implements ItemReader<Order> { private static final Map<String, Product> PRODUCT_CACHE = new ConcurrentHashMap<>(); @Override public Order read() throws Exception { // 从缓存读取产品信息... return new Order(...); } }这会导致严重问题:
- Spring Batch的Step执行完毕后,ItemReader Bean不会被销毁(Singleton Scope)
PRODUCT_CACHE持续增长,且无法被GC回收(强引用)- 多个Job并发执行时,缓存被所有Job共享,造成数据污染
正确做法:
- 使用
@StepScope让Reader随Step生命周期创建/销毁:
@Bean @StepScope public ItemReader<Order> orderReader(@Value("#{jobParameters['input.file']}") String file) { return new FlatFileItemReaderBuilder<Order>() .name("orderReader") .resource(new ClassPathResource(file)) .lineMapper(lineMapper()) .build(); }- 若必须缓存,使用
Caffeine并设置过期策略:
private final LoadingCache<String, Product> productCache = Caffeine.newBuilder() .maximumSize(10000) .expireAfterWrite(1, TimeUnit.HOURS) .build(key -> productService.findById(key));5. 监控与诊断:没有Metrics的Batch Job就像蒙眼开车
Spring Boot Actuator + Micrometer是标配,但Spring Batch的指标需要针对性埋点。默认情况下,Actuator只暴露/actuator/metrics中的jvm.*、http.server.requests等通用指标,对Batch Job的深度监控几乎为零。
5.1 必须暴露的5类核心指标
| 指标类别 | Micrometer名称 | 采集方式 | 业务意义 |
|---|---|---|---|
| Job执行成功率 | spring.batch.job.executions | JobExecutionListener中afterJob()上报 | 全局健康度,低于95%需告警 |
| Step处理吞吐量 | spring.batch.step.items.processed | StepExecutionListener中afterStep()上报 | 发现性能瓶颈(如某Step骤降) |
| Chunk失败率 | spring.batch.chunk.failures | ChunkListener中onError()上报 | 定位脏数据或临时故障 |
| 数据库连接等待 | jdbc.connections.active | HikariCP内置指标 | 判断是否需调大连接池 |
| JVM内存压力 | jvm.memory.used | Actuator默认 | 关联分析OOM根因 |
5.2 自定义ChunkListener实现细粒度监控
@Component public class BatchMetricsChunkListener implements ChunkListener { private final MeterRegistry meterRegistry; public BatchMetricsChunkListener(MeterRegistry meterRegistry) { this.meterRegistry = meterRegistry; } @Override public void beforeChunk(ChunkContext context) { // 记录Chunk开始时间 context.getStepContext().getStepExecution() .getExecutionContext() .putLong("chunk.start.time", System.currentTimeMillis()); } @Override public void afterChunk(ChunkContext context) { StepExecution stepExecution = context.getStepContext().getStepExecution(); long startTime = stepExecution.getExecutionContext().getLong("chunk.start.time"); long duration = System.currentTimeMillis() - startTime; // 上报处理条数 Counter.builder("spring.batch.chunk.items.processed") .tag("step", stepExecution.getStepName()) .tag("status", "success") .register(meterRegistry) .increment(stepExecution.getWriteCount()); // 上报耗时(毫秒) Timer.builder("spring.batch.chunk.duration") .tag("step", stepExecution.getStepName()) .register(meterRegistry) .record(duration, TimeUnit.MILLISECONDS); } @Override public void onError(ChunkContext context, Exception e) { // 上报失败数 Counter.builder("spring.batch.chunk.failures") .tag("step", context.getStepContext().getStepExecution().getStepName()) .tag("exception", e.getClass().getSimpleName()) .register(meterRegistry) .increment(); } }5.3 生产环境必备的3个诊断命令
当Job卡住时,不要急着重启,先执行以下诊断:
诊断1:查看活跃线程与锁
# 进入JVM进程 jstack -l <pid> | grep -A 20 "batch-" # 输出示例: "batch-1" #25 daemon prio=5 os_prio=0 cpu=12345.67ms elapsed=678.90s tid=0x00007f8b4c0a1000 nid=0x1a waiting for monitor entry [0x00007f8b3d7f9000] java.lang.Thread.State: BLOCKED (on object monitor) at com.zaxxer.hikari.pool.HikariPool.getConnection(HikariPool.java:188) - waiting to lock <0x000000071a2b3c80> (a com.zaxxer.hikari.pool.HikariPool)说明:线程在等待HikariCP连接,检查spring.datasource.hikari.maximum-pool-size是否过小。
诊断2:检查JobRepository状态
-- 查看最近10个失败Job SELECT JOB_INSTANCE_ID, JOB_NAME, START_TIME, END_TIME, STATUS, EXIT_CODE FROM BATCH_JOB_EXECUTION WHERE STATUS='FAILED' ORDER BY START_TIME DESC LIMIT 10; -- 查看某Step的详细上下文 SELECT STEP_NAME, READ_COUNT, WRITE_COUNT, STATUS, EXIT_CODE, SUBSTRING_INDEX(SUBSTRING_INDEX(EXT_CONTEXT, '"last.processed.id":"', -1), '"', 1) AS last_id FROM BATCH_STEP_EXECUTION WHERE JOB_EXECUTION_ID = 12345;诊断3:验证Chunk配置合理性
# 获取JVM内存使用详情 jstat -gc <pid> # 关键指标: # S0C/S1C:幸存者区容量(应远小于Chunk内存占用) # EC:伊甸园区容量(Chunk对象主要分配在此) # 如果EC持续100%,说明Chunk过大或GC策略不当6. 从示例到生产:一个真实订单对账Job的完整实现
现在,我们把前述所有原则,整合成一个可直接部署的订单对账Job。它解决的核心问题是:每日凌晨同步第三方支付平台的交易明细,比对本地订单状态,自动修正差异订单。
6.1 项目结构与依赖
pom.xml关键依赖:
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-batch</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-jdbc</artifactId> </dependency> <dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <scope>runtime</scope> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-actuator</artifactId> </dependency> <dependency> <groupId>io.micrometer</groupId> <artifactId>micrometer-registry-prometheus</artifactId> </dependency> </dependencies>6.2 配置文件:分离开发与生产
application.yml:
spring: batch: initialize-schema: always # 仅开发环境自动建表 job: enabled: false # 禁用自动启动,由Controller触发 datasource: url: jdbc:mysql://localhost:3306/main_db?useSSL=false&serverTimezone=Asia/Shanghai username: root password: root datasource: batch: url: jdbc:mysql://localhost:3306/batch_db?useSSL=false&serverTimezone=Asia/Shanghai username: batch_user password: batch_pass jackson: date-format: yyyy-MM-dd HH:mm:ss time-zone: Asia/Shanghai # Batch专属配置 batch: chunk: size: 200 # 经压测确定 retry: max-attempts: 3 job: reconciliation: input-date: "#{systemProperties['run.date'] ?: '20231001'}" # 支持系统属性传参6.3 核心组件实现
Step 1:ItemReader —— 从支付平台API拉取数据
@Bean @StepScope public ItemReader<PaymentRecord> paymentApiReader( @Value("#{jobParameters['run.date']}") String runDate, RestTemplate restTemplate) { return new ListItemReader<>(fetchPaymentRecords(runDate, restTemplate)); } private List<PaymentRecord> fetchPaymentRecords(String runDate, RestTemplate restTemplate) { // 调用支付平台API,参数:date=runDate String url = "https://api.payment.com/v1/transactions?date=" + runDate; try { ResponseEntity<PaymentResponse> response = restTemplate.getForEntity(url, PaymentResponse.class); return response.getBody().getData(); } catch (Exception e) { throw new RuntimeException("Failed to fetch payment records for " + runDate, e); } }Step 2:ItemProcessor —— 业务规则校验与转换
@Bean @StepScope public ItemProcessor<PaymentRecord, OrderReconciliation> reconciliationProcessor( @Value("#{jobParameters['run.date']}") String runDate) { return paymentRecord -> { // 1. 幂等性校验:跳过已处理过的支付记录 if (reconciliationService.isProcessed(paymentRecord.getTradeNo())) { return null; // 返回null表示跳过 } // 2. 本地订单查询 Order order = orderService.findByTradeNo(paymentRecord.getTradeNo()); if (order == null) { log.warn("Order not found for tradeNo: {}", paymentRecord.getTradeNo()); return null; } // 3. 状态比对 OrderReconciliation result = new OrderReconciliation(); result.setOrderId(order.getId()); result.setPaymentStatus(paymentRecord.getStatus()); result.setLocalStatus(order.getStatus()); result.setNeedUpdate(!paymentRecord.getStatus().equals(order.getStatus())); return result; }; }Step 3:ItemWriter —— 原子化更新订单状态
@Bean @StepScope public ItemWriter<OrderReconciliation> reconciliationWriter( @Qualifier("batchTransactionManager") PlatformTransactionManager transactionManager) { return items -> { // 批量更新,但必须在同一个事务内 for (OrderReconciliation item : items) { if (item.isNeedUpdate()) { orderService.updateStatus(item.getOrderId(), item.getPaymentStatus()); } } // 此处隐式提交事务 }; }Step 4:Job配置 —— 集成所有组件
@Bean public Job orderReconciliationJob( Step orderReconciliationStep, JobExecutionListener jobExecutionListener) { return jobBuilderFactory.get("orderReconciliationJob") .start(orderReconciliationStep) .listener(jobExecutionListener) .incrementer(new RunIdIncrementer()) // 每次生成新JobInstance .build(); } @Bean public Step orderReconciliationStep( @Qualifier("paymentApiReader") ItemReader<PaymentRecord> reader, @Qualifier("reconciliationProcessor") ItemProcessor<PaymentRecord, OrderReconciliation> processor, @Qualifier("reconciliationWriter") ItemWriter<OrderReconciliation> writer, @Qualifier("batchTransactionManager") PlatformTransactionManager transactionManager) { return stepBuilderFactory.get("orderReconciliationStep") .<PaymentRecord, OrderReconciliation>chunk(200) .reader(reader) .processor(processor) .writer(writer) .faultTolerant() .skip(PaymentApiException.class) // 支付平台API异常跳过 .skipLimit(100) // 最多跳过100条 .retry(RetryableException.class) // 可重试异常 .retryLimit(3) .backOffPolicy(new ExponentialBackOffPolicy() {{ setInitialInterval(1000L); setMultiplier(2.0); setMaxInterval(10000L); }}) .transactionManager(transactionManager) .build(); }6.4 启动与调度:脱离@Scheduled的优雅方案
创建BatchController提供HTTP触发入口:
@RestController @RequestMapping("/batch") public class BatchController { @Autowired private JobLauncher jobLauncher; @Autowired private Job orderReconciliationJob; @PostMapping("/reconcile") public ResponseEntity<String> triggerReconciliation( @RequestParam String runDate, @RequestParam(defaultValue = "false") boolean async) { try { JobParameters params = new JobParametersBuilder() .addString("run.date", runDate) .addLong("time", System.currentTimeMillis()) .toJobParameters(); JobExecution execution; if (async) { execution = jobLauncher.run(orderReconciliationJob, params); return ResponseEntity.accepted().body("Job started asynchronously: " + execution.getId()); } else { execution = jobLauncher.run(orderReconciliationJob, params); return ResponseEntity.ok("Job completed: " + execution.getStatus()); } } catch (Exception e) { return ResponseEntity.badRequest().body("Job failed: " + e.getMessage()); } } }调度方案对比:
- ❌
@Scheduled(fixedRate = 3600000):无法传递参数,无法监控单次执行状态 - ✅Quartz集成:支持集群调度、失败重试、参数化触发
- ✅Kubernetes CronJob:云原生首选,资源隔离彻底
- ✅HTTP API + 外部调度器(如Airflow):最灵活,可观测性最强
我推荐Airflow,因其能天然串联多个Batch Job(如:先跑对账Job,成功后再触发报表生成Job),且提供Web UI实时查看DAG状态。
7. 最后分享一个血泪教训:OOM时别急着加-Xmx
去年双十一前,一个库存同步Job在压测时频繁OOM,运维同事第一反应是“加堆内存”,从2G加到4G,问题依旧。最后发现根因是:
- ItemReader使用
JdbcCursorItemReader,但SQL未加ORDER BY id - MySQL优化器对无序查询使用了
Using filesort,导致临时表撑爆磁盘 JdbcCursorItemReader的游标机制在无序结果集上失效,不断重复读取同一数据
诊断过程:
jstack发现线
