SpringBatch 5实现重试&跳过策略回调:失败存储与末尾重试方案
处理失败记录存储与流程末尾重试方案
一、拆分两类失败场景的处理逻辑
1. 内容无效(数据层面错误)
这类错误是数据本身不符合规则(比如字段格式错误、必填项缺失),重试也无法解决,需要捕获并存储失败记录,同时跳过该条数据保证主流程继续。
- 首先定义自定义异常标记无效数据:
public class InvalidDataException extends RuntimeException { private final EstablishmentEntity invalidItem; public InvalidDataException(String message, EstablishmentEntity invalidItem) { super(message); this.invalidItem = invalidItem; } public EstablishmentEntity getInvalidItem() { return invalidItem; } }
- 在
ItemProcessor中加入校验逻辑,不通过时抛出上述异常:
@Bean public ItemProcessor<EstablishmentEntity, EstablishmentEntity> processor() { return item -> { // 示例校验:假设name字段不能为空 if (item.getName() == null || item.getName().trim().isEmpty()) { throw new InvalidDataException("名称不能为空", item); } // 其他校验逻辑... return item; }; }
- 修改Step配置,添加跳过规则和
SkipListener来收集失败记录:
@Bean public Step step1( BatchReadListener listener, ItemSkipListener skipListener, JobRepository jobRepository, JpaTransactionManager transactionManager, RepositoryItemWriter<EstablishmentEntity> writer) throws IOException { log.warn("step() called"); return new StepBuilder("step1", jobRepository) .<EstablishmentEntity, EstablishmentEntity> chunk(10, transactionManager) .reader(reader()) .processor(processor()) .writer(writer) .listener(listener) .listener(skipListener) .faultTolerant() .retryLimit(128) .retry(PessimisticLockException.class) .retry(PessimisticLockingFailureException.class) .retry(CannotAcquireLockException.class) .retry(SQLException.class) .skip(InvalidDataException.class) // 跳过内容无效的异常 .skipLimit(Integer.MAX_VALUE) // 根据业务需求调整最大跳过数 .build(); }
- 实现
SkipListener存储失败记录:
@Component public class ItemSkipListener implements SkipListener<EstablishmentEntity, EstablishmentEntity> { private final FailedItemRepository failedItemRepository; public ItemSkipListener(FailedItemRepository failedItemRepository) { this.failedItemRepository = failedItemRepository; } @Override public void onSkipInProcess(EstablishmentEntity item, Throwable t) { if (t instanceof InvalidDataException) { FailedItemEntity failedItem = new FailedItemEntity(); failedItem.setOriginalData(JSON.toJSONString(item)); // 序列化原数据存储 failedItem.setErrorMessage(t.getMessage()); failedItem.setFailureType("INVALID_DATA"); failedItem.setStatus("UNPROCESSED"); failedItemRepository.save(failedItem); } } @Override public void onSkipInRead(Throwable t) { // 处理CSV读取时的格式错误行 FailedItemEntity failedItem = new FailedItemEntity(); failedItem.setErrorMessage(t.getMessage()); failedItem.setFailureType("READ_ERROR"); failedItem.setStatus("UNPROCESSED"); failedItemRepository.save(failedItem); } @Override public void onSkipInWrite(EstablishmentEntity item, Throwable t) { // 写入时的非重试类错误(比如数据违反唯一约束) if (!(t instanceof SQLException || t instanceof PessimisticLockException)) { FailedItemEntity failedItem = new FailedItemEntity(); failedItem.setOriginalData(JSON.toJSONString(item)); failedItem.setErrorMessage(t.getMessage()); failedItem.setFailureType("WRITE_ERROR"); failedItem.setStatus("UNPROCESSED"); failedItemRepository.save(failedItem); } } }
2. MySQL连接/锁异常(系统层面临时错误)
这类错误(连接超时、死锁)是临时的,你的Step已经配置了重试,但如果重试耗尽仍失败,需要存储这些记录留到流程末尾重试。
- 实现
RetryListener捕获重试耗尽的记录:
@Component public class ItemRetryListener extends RetryListenerSupport { private final FailedItemRepository failedItemRepository; private static final int MAX_RETRY = 128; // 和Step配置的retryLimit保持一致 public ItemRetryListener(FailedItemRepository failedItemRepository) { this.failedItemRepository = failedItemRepository; } @Override public <T, E extends Throwable> void onError(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) { // 检查是否达到重试上限 if (context.getRetryCount() == MAX_RETRY - 1) { EstablishmentEntity item = (EstablishmentEntity) context.getAttribute("currentItem"); if (item != null) { FailedItemEntity failedItem = new FailedItemEntity(); failedItem.setOriginalData(JSON.toJSONString(item)); failedItem.setErrorMessage(throwable.getMessage()); failedItem.setFailureType("RETRY_EXHAUSTED"); failedItem.setRetryCount(MAX_RETRY); failedItem.setStatus("UNPROCESSED"); failedItemRepository.save(failedItem); } } } }
- 修改Step配置,注册
RetryListener并在Processor中传递当前Item到RetryContext:
@Bean public Step step1( BatchReadListener listener, ItemSkipListener skipListener, ItemRetryListener retryListener, JobRepository jobRepository, JpaTransactionManager transactionManager, RepositoryItemWriter<EstablishmentEntity> writer) throws IOException { log.warn("step() called"); return new StepBuilder("step1", jobRepository) .<EstablishmentEntity, EstablishmentEntity> chunk(10, transactionManager) .reader(reader()) .processor((ItemProcessor<EstablishmentEntity, EstablishmentEntity>) item -> { // 将当前Item存入RetryContext,方便监听器获取 RetrySynchronizationManager.getContext().setAttribute("currentItem", item); return processor().process(item); }) .writer(writer) .listener(listener) .listener(skipListener) .listener(retryListener) .faultTolerant() .retryLimit(MAX_RETRY) .retry(PessimisticLockException.class) .retry(PessimisticLockingFailureException.class) .retry(CannotAcquireLockException.class) .retry(SQLException.class) .skip(InvalidDataException.class) .skipLimit(Integer.MAX_VALUE) .build(); }
二、流程末尾执行重试操作
通过多Step的Job实现:主Step完成正常导入后,执行重试Step处理之前存储的RETRY_EXHAUSTED类型记录。
1. 重试Step的读取器
读取需要重试的失败记录:
@Bean public ItemReader<FailedItemEntity> failedItemReader(FailedItemRepository failedItemRepository) { // 仅读取需要重试的记录 List<FailedItemEntity> retryItems = failedItemRepository.findByFailureTypeAndStatus("RETRY_EXHAUSTED", "UNPROCESSED"); return new ListItemReader<>(retryItems); }
2. 重试Step的处理器
将序列化的原数据反序列化为实体类:
@Bean public ItemProcessor<FailedItemEntity, EstablishmentEntity> failedItemProcessor() { return failedItem -> JSON.parseObject(failedItem.getOriginalData(), EstablishmentEntity.class); }
3. 重试Step与Job配置
@Bean public Step retryStep( ItemReader<FailedItemEntity> failedItemReader, ItemProcessor<FailedItemEntity, EstablishmentEntity> failedItemProcessor, RepositoryItemWriter<EstablishmentEntity> writer, RetryWriteListener retryWriteListener, JobRepository jobRepository, JpaTransactionManager transactionManager) { return new StepBuilder("retryStep", jobRepository) .<FailedItemEntity, EstablishmentEntity> chunk(10, transactionManager) .reader(failedItemReader) .processor(failedItemProcessor) .writer(writer) .listener(retryWriteListener) .faultTolerant() .retryLimit(10) // 重试Step设置较小的重试次数 .retry(PessimisticLockException.class) .retry(SQLException.class) .build(); } @Bean public Job importJob(Step step1, Step retryStep, JobRepository jobRepository) { return new JobBuilder("importJob", jobRepository) .start(step1) .next(retryStep) // 主Step执行完毕后自动执行重试Step .build(); }
4. 重试后的状态更新
在WriteListener中更新失败记录的状态:
@Component public class RetryWriteListener implements ItemWriteListener<EstablishmentEntity> { private final FailedItemRepository failedItemRepository; public RetryWriteListener(FailedItemRepository failedItemRepository) { this.failedItemRepository = failedItemRepository; } @Override public void afterWrite(List<? extends EstablishmentEntity> items) { items.forEach(item -> { FailedItemEntity failedItem = failedItemRepository.findByOriginalData(JSON.toJSONString(item)); if (failedItem != null) { failedItem.setStatus("SUCCESS"); failedItemRepository.save(failedItem); } }); } @Override public void onWriteError(Exception exception, List<? extends EstablishmentEntity> items) { items.forEach(item -> { FailedItemEntity failedItem = failedItemRepository.findByOriginalData(JSON.toJSONString(item)); if (failedItem != null) { failedItem.setRetryCount(failedItem.getRetryCount() + 10); // 累加重试次数 failedItem.setStatus("FAILED"); failedItemRepository.save(failedItem); } }); } }
三、关键注意点
- 失败记录表需包含:原数据(序列化存储)、错误信息、失败类型、重试次数、状态、创建时间等字段。
- 内容无效的记录无需重试,直接标记为
UNPROCESSED等待人工处理;重试耗尽的记录若再次失败,标记为FAILED,后续可通过定时任务二次重试或人工介入。 - 重试Step的读取逻辑可按需调整,比如按时间范围读取、限制单次重试数量等。
内容的提问来源于stack exchange,提问作者davidvera
相关产品推荐
相关产品推荐

