如何在两个多租户数据源间高效复制数据(替代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
相关产品推荐
相关产品推荐

