Spring Batch中如何捕获并重试ChunkListener抛出的自定义异常?
解决Spring Batch ChunkListener异常无法被重试捕获的问题
问题原因
使用FaultTolerantStepBuilder时,配置的ChunkListener会被自动包装为TerminateOnExceptionChunkListenerDelegate,这个代理类会将监听器抛出的所有异常包装为FatalStepExecutionException——这是Spring Batch定义的致命异常,不会触发重试机制,导致你配置的retry(CustomException.class)规则失效。
可行解决方案
方案1:将锁检查逻辑迁移到ItemReader的首次读取环节
ItemReader的方法处于Step容错机制的覆盖范围内,抛出的CustomException会被重试规则捕获。我们可以在reader的首次读取前执行锁检查:
public class MyItemReader implements ItemReader<I> { private boolean isFirstChunk = true; @Override public I read() throws Exception { if (isFirstChunk) { // Redis分布式锁检查逻辑 if (isLocked) { throw new CustomException(); } isFirstChunk = false; } // 原有业务读取逻辑 return ...; } @Override public void open(ExecutionContext executionContext) throws ItemStreamException { // 重置标记,确保Step重启或Chunk循环时重新检查锁 isFirstChunk = true; super.open(executionContext); } }
方案2:自定义ChunkListener代理类,避免包装自定义异常
创建自定义代理类,重写异常处理逻辑,仅对非自定义异常做致命包装:
public class CustomChunkListenerDelegate extends TerminateOnExceptionChunkListenerDelegate { public CustomChunkListenerDelegate(ChunkListener delegate) { super(delegate); } @Override protected void handleException(Exception e) throws StepExecutionException { // 自定义异常直接抛出,不包装为致命异常 if (e instanceof CustomException) { throw (CustomException) e; } // 其他异常保持原有致命包装逻辑 super.handleException(e); } }
然后在Step配置中使用这个代理类包装你的CustomChunkListener:
public Step MyBatchStep() { return stepBuilderFactory.get("refineCarStatusStep") .<I, O>chunk(CHUNK_SIZE) .reader(MyItemReader()) .processor(MyItemProcessor()) .writer(MyItemWriter()) .faultTolerant() .retry(CustomException.class) .retryLimit(3) .listener(new CustomChunkListenerDelegate(new CustomChunkListener())) .build(); }
方案3:在ChunkListener内部手动实现重试逻辑
使用Spring Retry的RetryTemplate在监听器内部处理重试,无需依赖Step的容错配置:
import org.springframework.retry.support.RetryTemplate; import org.springframework.retry.policy.SimpleRetryPolicy; import java.util.Collections; public class CustomChunkListener implements ChunkListener { private final RetryTemplate retryTemplate; public CustomChunkListener() { // 配置重试规则:最多重试3次,仅针对CustomException SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(3, Collections.singletonMap(CustomException.class, true)); this.retryTemplate = new RetryTemplate(); retryTemplate.setRetryPolicy(retryPolicy); } @Override public void beforeChunk(ChunkContext context) { retryTemplate.execute(context1 -> { // Redis分布式锁检查逻辑 if (isLocked) { throw new CustomException(); } return null; }); } @Override public void afterChunk(ChunkContext context) { unlock(); } @Override public void afterChunkError(ChunkContext context) { unlock(); } }
内容的提问来源于stack exchange,提问作者Tims
相关产品推荐
相关产品推荐

