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 '"([^"]+)"'));
- 给表B的
- 分批更新避免锁表:
直接全表更新会长时间锁表,改用按主键分段的分批更新:
每次调整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
相关产品推荐
相关产品推荐

