Spring Batch中JdbcCursorItemReader断连后的重试行为咨询
问题解答
核心问题:JdbcCursorItemReader重试时是否会重新执行SQL?
JdbcCursorItemReader的核心行为如下:
- 正常流程下,SQL仅执行一次:Step启动时,Reader会执行配置的SQL获取
ResultSet,之后逐行读取数据填充Chunk,直到ResultSet耗尽。 - 重试触发时的行为分两种场景:
- 读取阶段触发重试:如果在读取Chunk数据时发生连接中断、ResultSet失效等异常(比如你配置的
ConnectException、SQLTransientConnectionException),重试会重新初始化Reader——也就是重新执行SQL查询,重新获取ResultSet。此时如果源数据已经变更,新的ResultSet会包含最新数据,可能导致重复或丢失。 - 写入阶段触发重试:如果写入Chunk时发生异常,重试仅会重新执行当前Chunk的写入操作,不会重新执行读取SQL,因为当前Chunk的读取已经完成,数据已经在内存中。
- 读取阶段触发重试:如果在读取Chunk数据时发生连接中断、ResultSet失效等异常(比如你配置的
你的数据丢失问题分析
结合你的场景,数据丢失可能的原因包括:
- UNDO_RETENTION时长不足:Oracle的快照依赖UNDO表空间,若迁移时间超过15分钟,快照数据可能被覆盖,会抛出
ORA-01555: snapshot too old错误。你配置了SQLRecoverableException重试,但如果没捕获到该异常,可能是:- 异常未被正确归类到你配置的重试类型中;
- 日志级别或监听器配置问题,导致异常未被记录。
- 移除数据边界条件:原来的边界条件(比如按ID分段)可以保证读取的一致性,移除后如果迁移期间源数据发生变更,即使游标基于快照,一旦快照失效重新执行SQL,会读取到新数据,导致原快照中未读取的数据丢失。
- Chunk过大:50000条的Chunk size过大,一旦读取或写入失败,重试的成本更高,且长时间持有ResultSet会增加UNDO快照过期的风险。
建议解决方案
- 验证重试行为:故意制造连接中断(比如临时断开Oracle连接),观察RetryListener是否输出日志,同时检查数据库的SQL执行记录,确认是否重新执行了读取SQL。
- 恢复数据边界+固定快照时间:使用Oracle闪回查询固定读取的快照时间,避免UNDO_RETENTION限制,同时按ID分段读取,示例SQL:
select t.id, t.group from teams AS OF TIMESTAMP TO_TIMESTAMP('2024-05-20 10:00:00', 'YYYY-MM-DD HH24:MI:SS') where t.id between ? and ? - 调整重试异常范围:明确添加
ORA-01555对应的异常处理,同时将日志级别调整为DEBUG,查看读取阶段的详细日志:.retry(OracleSQLException.class) // 若使用Oracle驱动的特定异常类 - 保证写入幂等性:修改写入SQL为幂等操作,避免重复写入或丢失:
INSERT INTO "teams" ("id", "group") VALUES (:id, :group) ON CONFLICT ("id") DO UPDATE SET "group" = EXCLUDED."group" - 替换为分页读取:使用
JdbcPagingItemReader替代JdbcCursorItemReader,每次查询一个分段的数据,避免长时间持有ResultSet,重试时仅重新查询当前分段,更稳定。
你的代码格式化
public class Configuration{ public final String READ_TEAM_TO_DATE_SQL = "select t.id, t.group from teams"; public final String WRITE_TEAM_SQL = "INSERT INTO \"teams\" (\"id\", \"group\") VALUES (:id, :group)"; @Autowired @Qualifier("dataSource") private DataSource postgresqlDatasource; @Autowired @Qualifier("oracleDataSource") private DataSource oracleDataSource; @Autowired private JobRepository jobRepository; @Autowired private PlatformTransactionManager transactionManager; @Bean public Flow getTeamsFlow() { return new FlowBuilder<SimpleFlow>("getTeams") .start(getTeamsStep()) .on("COMPLETED").end() .on("FAILED").fail() .build(); } @Bean public Step getTeamsStep() { return new StepBuilder("getTeamsStep", jobRepository) .<Team, Team>chunk(50_000, transactionManager) .faultTolerant() .retry(ConnectException.class) .retry(SocketTimeoutException.class) .retry(SQLTransientConnectionException.class) .retry(SQLTransientException.class) .retry(SQLTimeoutException.class) .retry(SQLRecoverableException.class) .retryLimit(10) .listener(new ItemWriteListener<>() { @Override public void onWriteError(Exception ex, Chunk<? extends Team> items) { log.error("Error during writing team - chunk level: {} ", items.getItems().stream().map(Team::toString) .collect(Collectors.joining("|"))); log.error("Error during writing team - chunk level - exception: ", ex); ItemWriteListener.super.onWriteError(ex, items); } }) .listener(new ItemReadListener<>() { @Override public void onReadError(Exception ex) { log.error("Error during reading team - chunk level - exception: ", ex); ItemReadListener.super.onReadError(ex); } }) .listener(new RetryListener() { @Override public <T, E extends Throwable> void close(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) { log.error("Occurred error during last retry: ", throwable); RetryListener.super.close(context, callback, throwable); } @Override public <T, E extends Throwable> void onError(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) { log.error("Occurred error during retry: ", throwable); RetryListener.super.onError(context, callback, throwable); } }) .reader(getTeamReader()) .writer(getTeamWriter()) .build(); } @Bean @StepScope public JdbcCursorItemReader<Team> getTeamReader() { JdbcCursorItemReader<Team> reader = new JdbcCursorItemReader<>(); reader.setSql(READ_TEAM_TO_DATE_SQL); reader.setDataSource(oracleDataSource); reader.setRowMapper((ResultSet rs, int rowNum) -> Team.builder() .id(rs.getString("id")) .group(rs.getString("group")) .build()); return reader; } @Bean public JdbcBatchItemWriter<Team> getTeamWriter() { JdbcBatchItemWriter<Team> writer = new JdbcBatchItemWriter<>(); writer.setDataSource(postgresqlDatasource); writer.setSql(WRITE_TEAM_SQL); writer.setItemSqlParameterSourceProvider(new BeanPropertyItemSqlParameterSourceProvider<>()); return writer; } @Bean public SimpleAsyncTaskExecutor taskExecutor() { return new SimpleAsyncTaskExecutor(); } }
内容的提问来源于stack exchange,提问作者RafalQA
相关产品推荐
相关产品推荐

