如何用Pandas DataFrame更新PostgreSQL数据库表?
PostgreSQL二手车数据更新:Pandas to_sql无冲突更新替代方案
问题概述
我维护一套二手车爬虫系统:爬虫脚本抓取二手车网站数据存入PostgreSQL,另一脚本每小时运行一次,读取数据库中所有车辆的finn_code,重新抓取车辆状态并更新数据库,目标是分析车辆售出速度。
初始入库用Pandas DataFrame的to_sql实现,但更新时尝试两种方法均报错:
- 使用
on_conflict参数:提示参数不存在
with engine.connect() as con: updatedCars.to_sql( con = con, name = 'cars', if_exists = 'append', index_label = 'finn_code', on_conflict={'index': True, 'do_update': True, 'set_': updatedCars.columns.tolist()} )
- 使用
if_row_exists参数:同样提示参数不存在
with engine.connect() as con: updatedCars.to_sql( con = con, name = 'cars', if_exists = 'append', index_label = 'finn_code', if_row_exists='update' )
当前环境:Pandas 1.4.3,配合sqlalchemy、psycopg2。查阅Pandas文档确认,这两个参数在1.4.3版本中均未引入,需寻找替代方案。
当前实现代码
def fetchCarsFromDatabase(): connection = psycopg2.connect( host="localhost", database="finnscraper", user="user", password="password", ) cursor = connection.cursor() sql = 'SELECT finn_code FROM cars;' cursor.execute(sql) result = cursor.fetchall() finn_codes = [] for row in result: finn_code = ''.join(str(element) for element in row) finn_codes.append(finn_code) return finn_codes def updateCars(): finn_codes = fetchCarsFromDatabase() urls = [] for finn_code in finn_codes: url = 'https://www.finn.no/car/used/ad.html?finnkode=' + finn_code urls.append(url) allCars = getAllCars(urls) updatedCars = pd.DataFrame(allCars) updatedCars = updatedCars.rename(columns = { 'Finn kode': 'finn_code', 'Modell': 'model', 'Totalpris': 'total_price', 'Status': 'status', 'Farge': 'color', 'Hjuldrift': 'wheel_type', 'Effekt': 'effect', 'Sylindervolum': 'cylinder_volume', 'Vekt': 'weight', 'Antall seter': 'seats', 'Karosseri': 'body', 'Bilen står i': 'country_location', 'Reg.nr.': 'registration', '1. gang registrert': 'first_registration_date', 'Fargebeskrivelse': 'color_description', 'Interiørfarge': 'interior_color', 'Chassis nr. (VIN)': 'chassis_id', 'Pris eks omreg': 'price_ex_registration', 'Antall eiere': 'amount_of_owners', 'Sist EU-godkjent': 'last_eu_approved_date', 'Neste frist for EU-kontroll': 'next_eu_control_date' }) updatedCars = updatedCars.set_index('finn_code') engine = create_engine('postgresql://user:password@/finnscraper') with engine.connect() as con: updatedCars.to_sql( con = con, name = 'cars', if_exists = 'append', )
可行替代方案
方案1:SQLAlchemy Core执行UPSERT语句
利用PostgreSQL原生的ON CONFLICT语法,构造UPSERT语句逐条更新:
from sqlalchemy import text import pandas as pd from sqlalchemy import create_engine def updateCars(): # 省略前面的抓取、数据处理逻辑... updatedCars = updatedCars.set_index('finn_code') engine = create_engine('postgresql://user:password@localhost/finnscraper') with engine.begin() as conn: # 定义需要更新的字段,这里以status为例,可按需添加其他字段 upsert_sql = text(""" INSERT INTO cars (finn_code, status) VALUES (:finn_code, :status) ON CONFLICT (finn_code) DO UPDATE SET status = EXCLUDED.status; """) # 遍历DataFrame执行更新 for finn_code, row in updatedCars.iterrows(): conn.execute(upsert_sql, {'finn_code': finn_code, 'status': row['status']})
说明:需确保finn_code在数据库中是主键或唯一约束,否则冲突逻辑不生效。如需更新多个字段,只需扩展INSERT和SET部分的字段即可。
方案2:psycopg2批量UPSERT(高效处理大量数据)
使用psycopg2的execute_values实现批量UPSERT,性能优于逐条执行:
import psycopg2 from psycopg2.extras import execute_values import pandas as pd def updateCars(): # 省略前面的抓取、数据处理逻辑... updatedCars = updatedCars.set_index('finn_code') # 转换为元组列表:(finn_code, status, ...) records = [] for finn_code, row in updatedCars.iterrows(): # 按需选择要更新的字段,这里以status为例 records.append((finn_code, row['status'])) # 连接数据库 conn = psycopg2.connect( host="localhost", database="finnscraper", user="user", password="password", ) cur = conn.cursor() # 构造批量UPSERT语句 upsert_sql = """ INSERT INTO cars (finn_code, status) VALUES %s ON CONFLICT (finn_code) DO UPDATE SET status = EXCLUDED.status; """ # 批量执行 execute_values(cur, upsert_sql, records) conn.commit() cur.close() conn.close()
说明:execute_values会自动将批量数据打包成高效的SQL语句,适合处理大量车辆的更新操作。
方案3:升级Pandas版本(若允许)
if_row_exists参数在Pandas 2.0+版本中引入,支持'update'选项。若可以升级Pandas,可直接使用:
engine = create_engine('postgresql://user:password@/finnscraper') with engine.connect() as con: updatedCars.to_sql( con = con, name = 'cars', if_exists = 'append', index_label = 'finn_code', if_row_exists='update' )
注意:升级前需测试现有代码的兼容性,避免版本冲突。
优化建议
- 仅更新需要变更的字段(如
status),减少数据库IO开销 - 为
finn_code添加唯一约束或主键,确保冲突检测生效 - 添加日志记录,跟踪更新失败的车辆ID,便于排查问题
内容的提问来源于stack exchange,提问作者Snasegh
相关产品推荐
相关产品推荐

