监听器详解
本章定位:全面掌握 Spring Batch 的监听器体系——从 JobExecutionListener 到 ItemReadListener,覆盖生命周期监控、性能统计和异常告警。
定义与作用
Spring Batch 提供了丰富的监听器(Listener)接口,允许开发者在批处理生命周期的关键节点插入自定义逻辑。监听器分为 Job 级别和 Step 级别两大类。
飞翔科技架构师白歌的比喻:
"监听器是生产线的监控摄像头。JobExecutionListener 盯着整条线有没有开工/完工,StepExecutionListener 盯着每个工位,ChunkListener 盯着每个批次,ItemReadListener/ItemWriteListener 盯着每一件产品。"
核心原理
监听器调用时序
监听器矩阵
| 监听器接口 | 级别 | 回调方法 | 典型用途 |
|---|---|---|---|
JobExecutionListener | Job | beforeJob / afterJob | 作业开始/结束日志、资源清理、通知 |
StepExecutionListener | Step | beforeStep / afterStep | Step 统计、自定义 ExitStatus |
ChunkListener | Chunk | beforeChunk / afterChunk / afterChunkError | Chunk 耗时、进度百分比 |
ItemReadListener | Item | beforeRead / afterRead / onReadError | 读异常记录 |
ItemProcessListener | Item | beforeProcess / afterProcess / onProcessError | 处理耗时监控 |
ItemWriteListener | Item | beforeWrite / afterWrite / onWriteError | 写统计、失败数据保存 |
SkipListener | Item | onSkipInRead / onSkipInProcess / onSkipInWrite | 被跳过数据的日志/归档 |
RetryListener | Item | open / close / onError | 重试次数统计 |
完整示例
场景一:飞翔科技——全链路性能监控
@Component
public class PerformanceMonitor implements
JobExecutionListener, StepExecutionListener, ChunkListener {
private long jobStartTime;
private long stepStartTime;
// === Job 级别 ===
@Override
public void beforeJob(JobExecution jobExecution) {
jobStartTime = System.currentTimeMillis();
System.out.println("=== Job [" + jobExecution.getJobInstance().getJobName()
+ "] 开始 ===");
}
@Override
public void afterJob(JobExecution jobExecution) {
long duration = System.currentTimeMillis() - jobStartTime;
System.out.println("=== Job 完成 | 状态: " + jobExecution.getStatus()
+ " | 耗时: " + duration + "ms ===");
}
// === Step 级别 ===
@Override
public void beforeStep(StepExecution stepExecution) {
stepStartTime = System.currentTimeMillis();
}
@Override
public ExitStatus afterStep(StepExecution stepExecution) {
long duration = System.currentTimeMillis() - stepStartTime;
System.out.println(" Step [" + stepExecution.getStepName() + "]"
+ " | read=" + stepExecution.getReadCount()
+ " | write=" + stepExecution.getWriteCount()
+ " | skip=" + stepExecution.getSkipCount()
+ " | 耗时 " + duration + "ms");
return stepExecution.getExitStatus();
}
// === Chunk 级别 ===
@Override
public void afterChunk(ChunkContext context) {
StepExecution stepExec = context.getStepContext().getStepExecution();
System.out.println(" Chunk #" + stepExec.getCommitCount()
+ " | 累计读取 " + stepExec.getReadCount());
}
}
运行输出:
=== Job [orderProcessingJob] 开始 ===
Chunk #1 | 累计读取 500
Chunk #2 | 累计读取 1000
Step [importStep] | read=5000 | write=4980 | skip=20 | 耗时 8500ms
=== Job 完成 | 状态: COMPLETED | 耗时: 9200ms ===
场景二:SkipListener 记录被跳过的数据
@Component
public class SkipDataRecorder implements SkipListener<OrderDTO, OrderEntity> {
@Override
public void onSkipInRead(Throwable t) {
System.err.println("[SKIP-READ] 行解析失败: " + t.getMessage());
}
@Override
public void onSkipInProcess(OrderDTO item, Throwable t) {
// 将被跳过的数据写入错误文件
System.err.println("[SKIP-PROCESS] orderId=" + item.getOrderId()
+ " | reason=" + t.getMessage());
}
@Override
public void onSkipInWrite(OrderEntity item, Throwable t) {
System.err.println("[SKIP-WRITE] orderId=" + item.getOrderId()
+ " | reason=" + t.getMessage());
}
}
易错场景与避坑
反例一:afterStep 中修改 ExitStatus 导致流程异常
// ❌ 将 FAILED 改为 COMPLETED → Job 继续执行,数据可能不完整
@Override
public ExitStatus afterStep(StepExecution stepExecution) {
if (stepExecution.getStatus() == BatchStatus.FAILED) {
return ExitStatus.COMPLETED; // ❌ 掩盖错误!
}
return stepExecution.getExitStatus();
}
反例二:监听器中执行耗时操作
// ❌ afterChunk 中调外部 API → 拖慢每个 Chunk
@Override
public void afterChunk(ChunkContext ctx) {
slackClient.sendMessage("Chunk completed"); // 每次 500ms
// 100 个 Chunk → 额外耗时 50 秒
}
面试高频考点
Q1:afterStep 返回的 ExitStatus 和 Step 的 BatchStatus 有什么区别?
BatchStatus 是 Spring Batch 内部状态(COMPLETED/FAILED/STARTED),由框架决定。ExitStatus 是
afterStep返回的自定义退出码,用于 Job 的条件路由(on("CUSTOM_CODE").to(...))。
Q2:SkipListener 和重试的交互?
SkipListener 只在数据最终被跳过时触发(重试耗尽后)。
RetryListener在每次重试时触发。如果需要记录重试历史,用RetryListener;记录最终被丢弃的数据,用SkipListener。
上一章:命令行与 Web 启动下一章:重启与重试