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

Spring Batch监听器抛出异常时如何获取对应记录?

解决方案

针对JdbcCursorItemReader读异常时无法获取对应记录的问题,推荐以下两种高效实现方式,无需存储所有记录:

方案一:自定义RowMapper捕获映射异常并记录原始数据

JdbcCursorItemReader通过RowMapper将ResultSet映射为实体,异常通常发生在映射阶段。可以包装RowMapper,在映射失败时直接读取ResultSet中的原始数据并记录:

@Component
public class LoggingRowMapper<T> implements RowMapper<T> {

    private final RowMapper<T> delegate;
    private final Logger log = LoggerFactory.getLogger(getClass());

    public LoggingRowMapper(RowMapper<T> delegate) {
        this.delegate = delegate;
    }

    @Override
    public T mapRow(ResultSet rs, int rowNum) throws SQLException {
        try {
            return delegate.mapRow(rs, rowNum);
        } catch (SQLException | RuntimeException e) {
            // 拼接当前行的所有列名和值
            StringBuilder rowData = new StringBuilder();
            ResultSetMetaData metaData = rs.getMetaData();
            int columnCount = metaData.getColumnCount();
            for (int i = 1; i <= columnCount; i++) {
                String columnName = metaData.getColumnName(i);
                Object value = rs.getObject(i);
                rowData.append(columnName).append("=").append(value).append(", ");
            }
            // 记录异常和原始数据
            log.error("读取第{}行数据失败,原始数据:{},异常详情:", rowNum, rowData.toString(), e);
            // 重新抛出异常,不影响Spring Batch的异常处理流程
            throw e;
        }
    }
}

然后在配置JdbcCursorItemReader时,用这个包装类替代原RowMapper:

@Bean
public JdbcCursorItemReader<SomeEntity> someEntityReader(DataSource dataSource) {
    RowMapper<SomeEntity> originalRowMapper = new BeanPropertyRowMapper<>(SomeEntity.class);
    LoggingRowMapper<SomeEntity> loggingRowMapper = new LoggingRowMapper<>(originalRowMapper);
    
    return new JdbcCursorItemReaderBuilder<SomeEntity>()
            .dataSource(dataSource)
            .sql("SELECT id, name, age FROM some_table")
            .rowMapper(loggingRowMapper)
            .name("someEntityReader")
            .build();
}

方案二:包装JdbcCursorItemReader在read方法中捕获异常

直接继承JdbcCursorItemReader,重写read方法,在捕获异常时读取当前ResultSet的原始数据:

public class LoggingJdbcCursorItemReader<T> extends JdbcCursorItemReader<T> {

    private final Logger log = LoggerFactory.getLogger(getClass());

    @Override
    public T read() throws Exception {
        try {
            return super.read();
        } catch (Exception e) {
            ResultSet rs = getCursor();
            if (rs != null && !rs.isClosed()) {
                StringBuilder rowData = new StringBuilder();
                ResultSetMetaData metaData = rs.getMetaData();
                int columnCount = metaData.getColumnCount();
                for (int i = 1; i <= columnCount; i++) {
                    String columnName = metaData.getColumnName(i);
                    Object value = rs.getObject(i);
                    rowData.append(columnName).append("=").append(value).append(", ");
                }
                log.error("读取数据失败,原始数据:{},异常详情:", rowData.toString(), e);
            }
            throw e;
        }
    }
}

配置时直接使用这个自定义Reader:

@Bean
public LoggingJdbcCursorItemReader<SomeEntity> someEntityReader(DataSource dataSource) {
    return new LoggingJdbcCursorItemReaderBuilder<SomeEntity>()
            .dataSource(dataSource)
            .sql("SELECT id, name, age FROM some_table")
            .rowMapper(new BeanPropertyRowMapper<>(SomeEntity.class))
            .name("someEntityReader")
            .build();
}

补充:结合SkipListener统一处理跳过逻辑

如果允许跳过异常记录,可以配合Spring Batch的容错机制和SkipListener,统一记录跳过的情况(需结合上述方案记录原始数据):

@Component
class SomeSkipListener implements SkipListener<SomeEntity, SomeEntity> {

    private final Logger log = LoggerFactory.getLogger(getClass());

    @Override
    public void onSkipInRead(Throwable t) {
        log.error("读取数据时触发跳过,异常类型:{},信息:{}", t.getClass().getName(), t.getMessage());
    }

    @Override
    public void onSkipInWrite(SomeEntity item, Throwable t) {}

    @Override
    public void onSkipInProcess(SomeEntity item, Throwable t) {}
}

在Step中配置跳过策略:

@Bean
public Step someStep(JobRepository jobRepository, PlatformTransactionManager transactionManager,
                     LoggingJdbcCursorItemReader<SomeEntity> reader,
                     SomeProcessor processor,
                     SomeWriter writer,
                     SomeSkipListener skipListener) {
    return new StepBuilder("someStep", jobRepository)
            .<SomeEntity, SomeEntity>chunk(10, transactionManager)
            .reader(reader)
            .processor(processor)
            .writer(writer)
            .faultTolerant()
            .skip(SQLException.class) // 指定需要跳过的异常类型
            .skipLimit(10) // 允许跳过的最大次数
            .listener(skipListener)
            .build();
}

这两种方案都只在异常发生时处理对应行的数据,不会存储所有记录,完全适配大数据量、异常极少的场景。

内容的提问来源于stack exchange,提问作者Akane

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 14:32:13