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

基于Python实现两个PostgreSQL数据库间的数据传输与同步

PostgreSQL跨库表数据同步的Python实现

同步需求

  • 若table_2中的记录在table_1中不存在,则删除该记录;
  • 若table_2中的记录与table_1对应记录存在差异,则仅更新差异部分;
  • 若table_1中的记录在table_2中不存在,则插入该记录。

SQL实现示例(同库场景)

CREATE TABLE table_1
(
    id   INT,
    name VARCHAR(50)
);

INSERT INTO table_1
VALUES (1, 'A'),
       (2, 'B'),
       (3, 'C'),
       (4, 'D'),
       (5, 'E'),
       (6, 'F');

CREATE TABLE table_2 AS
SELECT *
FROM table_1;

-- 模拟源表数据变更
INSERT INTO table_1
VALUES (7, 'Z');

DELETE
FROM table_1
WHERE name = 'C';

UPDATE table_1
SET name = 'W'
WHERE name = 'E';

-- 同步逻辑
DELETE
FROM table_2
WHERE NOT EXISTS (SELECT 1
                  FROM table_1
                  WHERE table_1.id = table_2.id);

UPDATE table_2
SET name = table_1.name
FROM table_1
WHERE table_2.id = table_1.id
  AND table_2.name != table_1.name;

INSERT INTO table_2 (id, name)
SELECT id, name
FROM table_1
WHERE NOT EXISTS (SELECT 1
                  FROM table_2
                  WHERE table_2.id = table_1.id);

Python跨库同步实现代码

import psycopg2
import psycopg2.extras  # 导入批量执行工具

# 连接源数据库
source_conn = psycopg2.connect(
    host="source_host",
    database="source_db",
    user="source_user",
    password="source_password"
)

# 连接目标数据库
target_conn = psycopg2.connect(
    host="target_host",
    database="target_db",
    user="target_user",
    password="target_password"
)

# 创建游标
source_cursor = source_conn.cursor()
target_cursor = target_conn.cursor()

# 查询源表table_1数据
select_query = "SELECT * FROM table_1;"
source_cursor.execute(select_query)
source_data_rows = source_cursor.fetchall()

# 查询目标表table_2数据
select_query = "SELECT * FROM table_2;"
target_cursor.execute(select_query)
target_data_rows = target_cursor.fetchall()

# 将源表数据转为字典,以id为键,方便快速查找
source_data = {row[0]: row[1:] for row in source_data_rows}

print("源表数据:")
for row in source_data_rows:
    print(row)

print("目标表数据:")
for row in target_data_rows:
    print(row)

# 初始化统计计数器
deleted_count = 0
updated_count = 0
inserted_count = 0

# 处理目标表数据:删除不存在于源表的记录,更新差异记录
for row in target_data_rows:
    target_id = row[0]
    target_data = row[1:]

    if target_id not in source_data:
        # 删除目标表中源表不存在的记录
        delete_query = "DELETE FROM table_2 WHERE id = %s;"
        target_cursor.execute(delete_query, (target_id,))
        deleted_count += 1
        print(f"删除ID为{target_id}的记录")
    else:
        source_data_row = source_data[target_id]
        if target_data != source_data_row:
            # 更新目标表中与源表有差异的记录
            update_query = "UPDATE table_2 SET name = %s WHERE id = %s;"
            target_cursor.execute(update_query, (*source_data_row, target_id))
            updated_count += 1
            print(f"更新ID为{target_id}的记录")

# 处理源表新增记录:插入到目标表
# 先收集目标表已存在的id集合
target_existing_ids = {row[0] for row in target_data_rows}
insert_values = [(id, *data) for id, data in source_data.items() if id not in target_existing_ids]

if insert_values:
    insert_query = "INSERT INTO table_2 (id, name) VALUES %s;"
    psycopg2.extras.execute_values(target_cursor, insert_query, insert_values)
    inserted_count = len(insert_values)
    print(f"插入{inserted_count}条新记录")

# 提交事务
target_conn.commit()

# 关闭连接和游标
source_cursor.close()
source_conn.close()

target_cursor.close()
target_conn.close()

# 输出同步结果
print(f"共删除{deleted_count}条记录")
print(f"共更新{updated_count}条记录")
print(f"共插入{inserted_count}条记录")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 12:37:02