基于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
相关产品推荐
相关产品推荐

