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

Spring Batch中JdbcCursorItemReader断连后的重试行为咨询

问题解答

核心问题:JdbcCursorItemReader重试时是否会重新执行SQL?

JdbcCursorItemReader的核心行为如下:

  • 正常流程下,SQL仅执行一次:Step启动时,Reader会执行配置的SQL获取ResultSet,之后逐行读取数据填充Chunk,直到ResultSet耗尽。
  • 重试触发时的行为分两种场景:
    1. 读取阶段触发重试:如果在读取Chunk数据时发生连接中断、ResultSet失效等异常(比如你配置的ConnectException、SQLTransientConnectionException),重试会重新初始化Reader——也就是重新执行SQL查询,重新获取ResultSet。此时如果源数据已经变更,新的ResultSet会包含最新数据,可能导致重复或丢失。
    2. 写入阶段触发重试:如果写入Chunk时发生异常,重试仅会重新执行当前Chunk的写入操作,不会重新执行读取SQL,因为当前Chunk的读取已经完成,数据已经在内存中。

你的数据丢失问题分析

结合你的场景,数据丢失可能的原因包括:

  • UNDO_RETENTION时长不足:Oracle的快照依赖UNDO表空间,若迁移时间超过15分钟,快照数据可能被覆盖,会抛出ORA-01555: snapshot too old错误。你配置了SQLRecoverableException重试,但如果没捕获到该异常,可能是:
    • 异常未被正确归类到你配置的重试类型中;
    • 日志级别或监听器配置问题,导致异常未被记录。
  • 移除数据边界条件:原来的边界条件(比如按ID分段)可以保证读取的一致性,移除后如果迁移期间源数据发生变更,即使游标基于快照,一旦快照失效重新执行SQL,会读取到新数据,导致原快照中未读取的数据丢失。
  • Chunk过大:50000条的Chunk size过大,一旦读取或写入失败,重试的成本更高,且长时间持有ResultSet会增加UNDO快照过期的风险。

建议解决方案

  1. 验证重试行为:故意制造连接中断(比如临时断开Oracle连接),观察RetryListener是否输出日志,同时检查数据库的SQL执行记录,确认是否重新执行了读取SQL。
  2. 恢复数据边界+固定快照时间:使用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 ?
    
  3. 调整重试异常范围:明确添加ORA-01555对应的异常处理,同时将日志级别调整为DEBUG,查看读取阶段的详细日志:
    .retry(OracleSQLException.class) // 若使用Oracle驱动的特定异常类
    
  4. 保证写入幂等性:修改写入SQL为幂等操作,避免重复写入或丢失:
    INSERT INTO "teams" ("id", "group") VALUES (:id, :group)
    ON CONFLICT ("id") DO UPDATE SET "group" = EXCLUDED."group"
    
  5. 替换为分页读取:使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 08:20:55