You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 12:05:12