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

PostgreSQL千万级数据同步SQL方案及Python分批处理咨询

PostgreSQL千万级数据同步方案

1. 核心SQL同步语句

首先需要从表A的IdTokenStore字段中提取出与表BIdTokenImport匹配的token字符串。观察IdTokenStore格式为(details, "{\"token-string-xxx\"}"),可通过正则表达式提取核心token:

UPDATE table_a a
SET CodeIntegerA = b.CodeIntegerB
FROM table_b b
WHERE substring(a.IdTokenStore from '"([^"]+)"') = b.IdTokenImport;

若正则效率偏低,也可使用字符串截取组合(适配固定格式):

UPDATE table_a a
SET CodeIntegerA = b.CodeIntegerB
FROM table_b b
WHERE trim(both '"' from split_part(split_part(a.IdTokenStore, ',', 2), '{', 2)) = b.IdTokenImport;

2. 千万级数据量优化方案

  • 添加针对性索引:
    • 给表B的IdTokenImport建唯一索引(若token唯一):CREATE UNIQUE INDEX idx_b_idtokenimport ON table_b(IdTokenImport);
    • 给表A的提取后token建表达式索引:CREATE INDEX idx_a_extracted_token ON table_a(substring(IdTokenStore from '"([^"]+)"'));
  • 分批更新避免锁表:
    直接全表更新会长时间锁表,改用按主键分段的分批更新:
    WITH batch AS (
        SELECT IdA FROM table_a 
        WHERE CodeIntegerA IS NULL 
        LIMIT 500000 OFFSET 0
    )
    UPDATE table_a a
    SET CodeIntegerA = b.CodeIntegerB
    FROM table_b b, batch
    WHERE a.IdA = batch.IdA
      AND substring(a.IdTokenStore from '"([^"]+)"') = b.IdTokenImport;
    
    每次调整OFFSET值循环执行即可。
  • 临时禁用非必要约束/触发器:若表A有更新触发器或非核心约束,临时禁用可大幅提升更新速度,完成后再恢复。
  • 更新统计信息:执行ANALYZE table_a; ANALYZE table_b;让优化器生成更优执行计划。
  • 开启并行查询:调整max_parallel_workers_per_gather参数,允许关联查询使用并行计算。

3. Python(Jupyter Notebook)分批处理实现

使用psycopg2库连接数据库,循环分批处理,示例代码如下:

import psycopg2
from psycopg2.extras import execute_batch

# 数据库连接参数
conn_params = {
    "dbname": "你的数据库名",
    "user": "你的用户名",
    "password": "你的密码",
    "host": "你的数据库地址"
}

batch_size = 500000  # 每次处理50万条
offset = 0

conn = psycopg2.connect(**conn_params)
cur = conn.cursor()

try:
    while True:
        # 获取当前批次的待更新IdA列表
        cur.execute("""
            SELECT IdA FROM table_a 
            WHERE CodeIntegerA IS NULL 
            LIMIT %s OFFSET %s
        """, (batch_size, offset))
        ids = cur.fetchall()
        
        if not ids:
            break  # 无剩余数据,退出循环
        
        # 批量执行更新
        update_sql = """
            UPDATE table_a a
            SET CodeIntegerA = b.CodeIntegerB
            FROM table_b b
            WHERE a.IdA = %s
              AND substring(a.IdTokenStore from '"([^"]+)"') = b.IdTokenImport
        """
        execute_batch(cur, update_sql, ids)
        
        conn.commit()
        offset += batch_size
        print(f"已完成 {offset} 条记录同步")
finally:
    cur.close()
    conn.close()

注意事项:

  • 可根据数据库负载调整batch_size,避免单次更新占用过多资源。
  • 若token提取逻辑复杂,建议先将表A的提取后token存入临时表,再关联更新,减少重复计算。
  • 执行过程中需监控数据库CPU、IO负载,避免影响其他业务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 19:13:12