如何用同列新DataFrame更新旧数据并同步至SQL数据库?
解决方案:用PostgreSQL更新已有表的对应行
你当前使用if_exists='append'的方式只会将new_data的行追加到表中,无法实现已有行的更新。要达成更新需求,需利用PostgreSQL的INSERT ... ON CONFLICT DO UPDATE语法,核心是先确定唯一匹配键——这里以Timestamp+Symbols作为联合主键,因为同一时间同一交易所的记录具备唯一性。
步骤1:给数据库表添加联合主键约束
如果你的表尚未设置主键,先执行以下SQL创建约束:
ALTER TABLE "table" ADD CONSTRAINT table_pkey PRIMARY KEY ("Timestamp", "Symbols");
也可通过SQLAlchemy执行:
from sqlalchemy import create_engine engine = create_engine('postgresql+psycopg2://xx:xxxx@localhost/stablecoin_db') with engine.connect() as conn: conn.execute('ALTER TABLE "table" ADD CONSTRAINT table_pkey PRIMARY KEY ("Timestamp", "Symbols");') conn.commit()
方法一:临时表+ON CONFLICT UPDATE(推荐,适配中等数据量)
先将new_data导入临时表,再执行更新语句:
import pandas as pd from sqlalchemy import create_engine engine = create_engine('postgresql+psycopg2://xx:xxxx@localhost/stablecoin_db') # 1. 将new_data导入临时表 new_data.to_sql('temp_table', engine, if_exists='replace', index=False) # 2. 执行更新:主键冲突时,更新USDT和USDC字段 with engine.connect() as conn: update_query = ''' INSERT INTO "table" ("Timestamp", "Symbols", "USDT", "USDC") SELECT "Timestamp", "Symbols", "USDT", "USDC" FROM temp_table ON CONFLICT ("Timestamp", "Symbols") DO UPDATE SET "USDT" = EXCLUDED."USDT", "USDC" = EXCLUDED."USDC"; ''' conn.execute(update_query) conn.commit() # 可选:删除临时表 with engine.connect() as conn: conn.execute('DROP TABLE temp_table;') conn.commit()
方法二:逐行更新(仅适配小数据量)
直接遍历new_data的每一行执行UPDATE语句:
import pandas as pd from sqlalchemy import create_engine engine = create_engine('postgresql+psycopg2://xx:xxxx@localhost/stablecoin_db') with engine.connect() as conn: for _, row in new_data.iterrows(): update_query = ''' UPDATE "table" SET "USDT" = %(usdt)s, "USDC" = %(usdc)s WHERE "Timestamp" = %(ts)s AND "Symbols" = %(sym)s; ''' conn.execute(update_query, { 'ts': row['Timestamp'], 'sym': row['Symbols'], 'usdt': row['USDT'], 'usdc': row['USDC'] }) conn.commit()
注意:该方法效率较低,数据量大时不建议使用。
方法三:psycopg2 copy_from+ON CONFLICT(适配大数据量,效率最优)
利用psycopg2的copy_from快速导入数据,再执行更新:
import pandas as pd import psycopg2 from io import StringIO # 创建数据库连接 conn = psycopg2.connect('postgresql://xx:xxxx@localhost/stablecoin_db') cur = conn.cursor() # 将new_data转为CSV格式字符串 output = StringIO() new_data.to_csv(output, sep='\t', header=False, index=False) output.seek(0) # 导入临时表 cur.execute('CREATE TEMP TABLE temp_table AS SELECT * FROM "table" LIMIT 0;') cur.copy_from(output, 'temp_table', sep='\t', columns=new_data.columns) # 执行更新 cur.execute(''' INSERT INTO "table" ("Timestamp", "Symbols", "USDT", "USDC") SELECT "Timestamp", "Symbols", "USDT", "USDC" FROM temp_table ON CONFLICT ("Timestamp", "Symbols") DO UPDATE SET "USDT" = EXCLUDED."USDT", "USDC" = EXCLUDED."USDC"; ''') conn.commit() cur.close() conn.close()
关键说明
- 必须保证
Timestamp和Symbols的组合在表中唯一,否则主键约束会报错,也无法正确匹配更新行。 - 若表字段有变动,需调整SQL语句中的字段名,确保与DataFrame列名一致。
内容的提问来源于stack exchange,提问作者Đỗ Văn Thắng
相关产品推荐
相关产品推荐

