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

使用Clojure & JDBC迁移5,000,000行数据至另一PostgreSQL数据库

兄弟,我太懂你这种踩坑的感受了——500万行数据,那些只会教你一次性把结果集全塞进内存的教程完全是站着说话不腰疼!咱们得换个思路,用分批次游标读取+分批次事务写入的方案,既不炸内存,又能保证迁移效率和数据安全。

核心思路拆解

本质就是避免一次性加载全量数据到内存,而是用「游标式分页读取」+「小批次事务写入」的组合拳,同时配合HikariCP和PostgreSQL的特性优化性能。

1. 用「键值分页」代替LIMIT+OFFSET(避免性能雪崩)

直接用OFFSET处理百万级数据会越跑越慢,因为数据库得先扫描前面所有行再返回结果。换成基于唯一有序字段(比如自增ID、时间戳)的键值分页,每次只取「上一批最后一条数据之后的N行」:

-- 示例:按自增ID分页,每次取1000行
SELECT id, col1, col2 FROM source_table WHERE id > ? ORDER BY id LIMIT 1000;

如果没有自增ID,就用唯一且有序的字段组合(比如created_at+id),避免因重复值漏数据。

2. 开启PostgreSQL的服务器端游标(关键!)

JDBC默认会把查询结果一次性拉到本地内存,哪怕你用了分页也没用。必须给查询设置fetchSize = Integer.MIN_VALUE,触发PostgreSQL的服务器端游标,让数据按需流式返回。

3. 小批次事务提交,别搞大事务

把整个迁移放一个大事务会锁死目标表,还容易因网络/连接问题前功尽弃。建议每1000~5000行提交一次事务,失败了也能从上次的最后ID断点续传。

代码示例(Java + HikariCP)
import com.zaxxer.hikari.HikariDataSource;
import java.sql.*;

public class PostgresDataMigration {
    public static void main(String[] args) {
        // 初始化源库和目标库的Hikari连接池(提前配置好参数)
        HikariDataSource sourceHikari = getHikariDataSource("source-db-config");
        HikariDataSource targetHikari = getHikariDataSource("target-db-config");

        long lastId = 0;
        int batchSize = 1000;
        boolean hasMoreData = true;

        try (Connection sourceConn = sourceHikari.getConnection();
             Connection targetConn = targetHikari.getConnection()) {

            targetConn.setAutoCommit(false); // 关闭自动提交,手动控制批次

            // 源查询:键值分页+开启服务器端游标
            String selectSql = "SELECT id, col1, col2, col3 FROM source_table WHERE id > ? ORDER BY id LIMIT ?";
            try (PreparedStatement selectStmt = sourceConn.prepareStatement(selectSql)) {
                selectStmt.setFetchSize(Integer.MIN_VALUE); // PostgreSQL游标触发关键配置

                String insertSql = "INSERT INTO target_table (id, col1, col2, col3) VALUES (?, ?, ?, ?)";
                try (PreparedStatement insertStmt = targetConn.prepareStatement(insertSql)) {

                    while (hasMoreData) {
                        selectStmt.setLong(1, lastId);
                        selectStmt.setInt(2, batchSize);

                        try (ResultSet rs = selectStmt.executeQuery()) {
                            int currentBatchCount = 0;
                            while (rs.next()) {
                                // 填充插入参数
                                insertStmt.setLong(1, rs.getLong("id"));
                                insertStmt.setString(2, rs.getString("col1"));
                                insertStmt.setInt(3, rs.getInt("col2"));
                                insertStmt.setTimestamp(4, rs.getTimestamp("col3"));
                                insertStmt.addBatch();
                                currentBatchCount++;
                                lastId = rs.getLong("id"); // 更新最后一条数据的ID
                            }

                            if (currentBatchCount == 0) {
                                hasMoreData = false; // 没有更多数据,结束循环
                                break;
                            }

                            // 执行批次插入并提交
                            insertStmt.executeBatch();
                            targetConn.commit();
                            System.out.printf("已迁移 %d 行,最后处理ID: %d%n", currentBatchCount, lastId);
                        } catch (SQLException e) {
                            targetConn.rollback(); // 出错回滚当前批次
                            throw new RuntimeException("批次迁移失败,断点ID: " + lastId, e);
                        }
                    }
                }
            }
            System.out.println("全量数据迁移完成!");
        } catch (SQLException e) {
            throw new RuntimeException("迁移过程异常", e);
        } finally {
            sourceHikari.close();
            targetHikari.close();
        }
    }

    // 封装Hikari连接池初始化(根据你的配置调整)
    private static HikariDataSource getHikariDataSource(String dbConfigKey) {
        HikariDataSource ds = new HikariDataSource();
        ds.setJdbcUrl("jdbc:postgresql://host:port/dbname");
        ds.setUsername("username");
        ds.setPassword("password");
        ds.setMaximumPoolSize(10); // 根据服务器配置调整连接数
        ds.setConnectionTimeout(30000);
        return ds;
    }
}
针对Clojure的适配方案(jdbc/with-db-transaction)

如果你用的是Clojure的jdbc库,核心逻辑一样,重点是设置:fetch-size启用游标,并且用小批次事务:

(require '[clojure.java.jdbc :as jdbc])

;; 提前初始化好源库和目标库的Hikari数据源
(def source-db {:datasource source-hikari})
(def target-db {:datasource target-hikari})

(def batch-size 1000)
(def last-id (atom 0))

(defn migrate-data []
  (loop []
    (let [rows (jdbc/query source-db
                           ["SELECT id, col1, col2, col3 FROM source_table WHERE id > ? ORDER BY id LIMIT ?" @last-id batch-size]
                           {:fetch-size Integer/MIN_VALUE})] ; 关键:启用PostgreSQL游标
      (if (empty? rows)
        (println "✅ 数据迁移完成!")
        (do
          ;; 小批次事务写入
          (jdbc/with-db-transaction [tx target-db]
            (jdbc/insert-multi! tx :target_table rows))
          ;; 更新断点ID
          (reset! last-id (-> rows last :id))
          (printf "已迁移 %d 行,当前断点ID: %d%n" (count rows) @last-id)
          (recur))))))

(migrate-data)
额外优化建议
  • 迁移前禁用目标表的索引和外键约束,迁移完成后再重建——否则每插入一行都要更新索引,速度会慢好几倍。
  • 调整HikariCP的maximumPoolSize:不要设太大(避免数据库连接过载),建议8~16之间,根据服务器CPU核心数调整。
  • 加个断点续传逻辑:把每次的lastId写到本地文件/数据库,万一迁移中断,下次启动直接从断点开始。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 04:08:23