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
相关产品推荐
相关产品推荐

