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

如何在两个多租户数据源间高效复制数据(替代ResultSet手动迭代)?

嘿,我太懂你现在的困扰了——手动遍历ResultSet逐行插入数据,数据量一大就会因为频繁的网络请求和事务提交变得巨慢。既然你两个都是PostgreSQL库,而且表结构完全一致,给你几个高效的替代方案,按性能优先级排序:

1. 用PostgreSQL原生COPY命令(性能天花板)

PostgreSQL的COPY命令是专门为批量数据导入/导出设计的,性能比普通INSERT高几个数量级,因为它直接绕过了常规的SQL解析流程,以二进制或文本流的方式传输数据。而且我们可以用管道流直接在两个库间传输,不用落地磁盘。

try (Connection mainConn = provider.getConnection("main");
     Connection altConn = provider.getConnection("alt")) {

    // 强转为PostgreSQL专属连接,获取CopyManager工具类
    PgConnection mainPgConn = mainConn.unwrap(PgConnection.class);
    PgConnection altPgConn = altConn.unwrap(PgConnection.class);
    CopyManager mainCopyManager = mainPgConn.getCopyAPI();
    CopyManager altCopyManager = altPgConn.getCopyAPI();

    // 用管道流实现两个库间的直接数据传输,无需中间存储
    PipedInputStream inStream = new PipedInputStream();
    PipedOutputStream outStream = new PipedOutputStream(inStream);

    // 异步执行主库数据导出,避免主线程阻塞
    new Thread(() -> {
        try {
            // 从主库导出数据到输出流(默认是二进制格式,效率最高)
            mainCopyManager.copyOut("COPY (SELECT * FROM companies) TO STDOUT", outStream);
        } catch (SQLException e) {
            e.printStackTrace();
        } finally {
            try {
                outStream.close();
            } catch (IOException e) {
                e.printStackTrace();
            }
        }
    }).start();

    // 从输入流导入数据到副库
    altCopyManager.copyIn("COPY companies FROM STDIN", inStream);
    altConn.commit();
} catch (SQLException | IOException e) {
    e.printStackTrace();
    // 异常时执行回滚
    try {
        if (altConn != null) altConn.rollback();
    } catch (SQLException ex) {
        ex.printStackTrace();
    }
}

2. JDBC批量插入优化(快速改造现有代码)

如果暂时没法用COPY,那先把你现有代码的两个致命问题解决:每次循环创建PreparedStatement和每次提交事务。优化后能大幅提升效率:

String getCompaniesQuery = "select * from companies";
String setRecordQuery = "insert into companies (company) values (?)";
// 批量大小:根据内存和网络情况调整,一般1000-5000条比较合适
int batchSize = 1000;
int recordCount = 0;

try (Connection main = provider.getConnection("main");
     Connection alt = provider.getConnection("alt");
     PreparedStatement queryStmt = main.prepareStatement(getCompaniesQuery);
     PreparedStatement insertStmt = alt.prepareStatement(setRecordQuery);
     ResultSet rs = queryStmt.executeQuery()) {

    // 关闭自动提交,手动控制事务提交时机
    alt.setAutoCommit(false);

    while (rs.next()) {
        insertStmt.setString(1, rs.getString(1));
        insertStmt.addBatch();
        recordCount++;

        // 达到批量阈值时执行一次批量插入并提交
        if (recordCount % batchSize == 0) {
            insertStmt.executeBatch();
            alt.commit();
            insertStmt.clearBatch(); // 清空批次,准备下一轮
        }
    }

    // 处理剩余的不足批量大小的数据
    if (recordCount % batchSize != 0) {
        insertStmt.executeBatch();
        alt.commit();
    }

} catch (SQLException e) {
    e.printStackTrace();
    // 异常时回滚事务
    try {
        if (alt != null) alt.rollback();
    } catch (SQLException ex) {
        ex.printStackTrace();
    }
}

3. 使用Spring Batch(适合复杂同步场景)

如果你的需求是定期同步数据、需要监控同步进度、异常重试或分片处理,Spring Batch是绝佳选择。它内置了批量处理的全套能力,而且适配JDBC数据源非常方便:

配置读取主库的ItemReader

@Bean
public JdbcCursorItemReader<Company> companyReader(@Qualifier("mainDataSource") DataSource mainDataSource) {
    return new JdbcCursorItemReaderBuilder<Company>()
            .dataSource(mainDataSource)
            .sql("SELECT * FROM companies")
            .rowMapper((rs, rowNum) -> new Company(rs.getString("company"))) // 映射为实体类
            .build();
}

配置写入副库的ItemWriter

@Bean
public JdbcBatchItemWriter<Company> companyWriter(@Qualifier("altDataSource") DataSource altDataSource) {
    return new JdbcBatchItemWriterBuilder<Company>()
            .dataSource(altDataSource)
            .sql("INSERT INTO companies (company) VALUES (:company)")
            .beanMapped() // 自动映射实体类属性到SQL参数
            .build();
}

配置Step和Job

@Bean
public Step syncCompaniesStep(ItemReader<Company> reader, 
                              ItemWriter<Company> writer,
                              JobRepository jobRepository,
                              PlatformTransactionManager transactionManager) {
    return new StepBuilder("syncCompaniesStep", jobRepository)
            .<Company, Company>chunk(1000, transactionManager) // 每1000条为一个批次
            .reader(reader)
            .writer(writer)
            .build();
}

@Bean
public Job syncCompaniesJob(Step syncCompaniesStep, JobRepository jobRepository) {
    return new JobBuilder("syncCompaniesJob", jobRepository)
            .start(syncCompaniesStep)
            .build();
}

运行这个Job就能自动完成批量同步,Spring Batch会帮你处理事务、异常、进度统计等细节。

4. MyBatis批量插入(如果项目已用MyBatis)

如果你的项目已经在用MyBatis,有两种高效批量插入方式:

方式一:用foreach生成批量INSERT语句

<!-- 在Mapper.xml中定义批量插入方法 -->
<insert id="batchInsertCompanies">
    INSERT INTO companies (company)
    VALUES
    <foreach collection="companyList" item="item" separator=",">
        (#{item.company})
    </foreach>
</insert>

方式二:用ExecutorType.BATCH模式

// 打开BATCH模式的SqlSession
try (SqlSession session = sqlSessionFactory.openSession(ExecutorType.BATCH)) {
    CompanyMapper mapper = session.getMapper(CompanyMapper.class);
    // 从主库查询所有数据
    List<Company> companies = mapper.selectAllCompanies();
    // 批量插入
    for (Company company : companies) {
        mapper.insertCompany(company);
    }
    session.commit(); // 统一提交
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:45:01