使用SQLAlchemy 2.x与PostgreSQL实现Upsert时数据库未更新的问题
问题:SQLAlchemy执行Upsert后数据库无更新
场景说明
通过Python调用外部API获取数据,转换后与PostgreSQL现有数据对比,生成包含新记录、变更记录的DataFrame,需要通过SQLAlchemy实现:
- 新记录直接插入
- 已有记录(按主键匹配)更新
现有实现代码
def update_absence(year): api_result = get_absence(year) db_result = get_database_absence(year) df = compare_dataframes(api_result, db_result, 'id') metadata_obj = MetaData() metadata_obj.reflect(bind=engine) some_table = Table("tb_absence", metadata_obj, autoload_with=engine) for item in df.to_dict('records'): insert_stmt = insert(some_table).values(item).on_conflict_do_update(constraint='tb_absence_pkey', set_=item) print(insert_stmt.compile()) with engine.connect() as conn: result = conn.execute(insert_stmt) print(result.rowcount) conn.commit()
问题现象
- 编译生成的SQL语句如下(符合Upsert逻辑):
INSERT INTO tb_absence (id, start_date, end_date, half_day, morning, user_id, employee_id, type, extra_vacation, state, substitute_state, workdays, hours, medical_certificate, comments, substitute_user_id, name) VALUES (%(id)s, %(start_date)s, %(end_date)s, %(half_day)s, %(morning)s, %(user_id)s, %(employee_id)s, %(type)s, %(extra_vacation)s, %(state)s, %(substitute_state)s, %(workdays)s, %(hours)s, %(medical_certificate)s, %(comments)s, %(substitute_user_id)s, %(name)s) ON CONFLICT ON CONSTRAINT tb_absence_pkey DO UPDATE SET id = %(param_1)s, start_date = %(param_2)s, end_date = %(param_3)s, half_day = %(param_4)s, morning = %(param_5)s, user_id = %(param_6)s, employee_id = %(param_7)s, type = %(param_8)s, extra_vacation = %(param_9)s, state = %(param_10)s, substitute_state = %(param_11)s, workdays = %(param_12)s, hours = %(param_13)s, medical_certificate = %(param_14)s, comments = %(param_15)s, substitute_user_id = %(param_16)s, name = %(param_17)s
- 循环中每条记录的
rowcount均返回1,但数据库始终无更新
问题原因
- 连接与事务管理错误:每次循环都创建新的
engine.connect()上下文,上下文退出时连接自动关闭,事务未提交。最后一行的conn.commit()引用的是最后一次循环的连接,此时该连接已经被关闭,提交操作无效。 - Upsert更新字段不规范:将主键
id包含在set_参数中,虽然PostgreSQL不会报错,但属于冗余操作,且可能导致SQLAlchemy参数命名混乱(如编译后的param_*)。
修复后的完整代码
from sqlalchemy import MetaData, Table, insert, engine def update_absence(year): api_result = get_absence(year) db_result = get_database_absence(year) df = compare_dataframes(api_result, db_result, 'id') if df.empty: print("无需要更新的记录") return metadata_obj = MetaData() metadata_obj.reflect(bind=engine) some_table = Table("tb_absence", metadata_obj, autoload_with=engine) # 统一创建连接,整个Upsert过程用同一个事务 with engine.connect() as conn: for item in df.to_dict('records'): # 构造Upsert语句:排除主键id,用excluded引用插入的字段值 update_dict = {col.name: some_table.c[col.name] for col in some_table.columns if col.name != 'id'} insert_stmt = insert(some_table).values(item).on_conflict_do_update( constraint='tb_absence_pkey', set_=update_dict ) # 或者更简洁的方式:使用excluded关键字引用插入的字段 # insert_stmt = insert(some_table).values(item).on_conflict_do_update( # constraint='tb_absence_pkey', # set_={col: insert_stmt.excluded[col] for col in some_table.columns if col.name != 'id'} # ) result = conn.execute(insert_stmt) print(f"处理记录ID {item['id']},影响行数:{result.rowcount}") # 统一提交事务 conn.commit() print("所有记录处理完成")
关键提示
- 事务范围:将所有Upsert操作放在同一个连接上下文内,确保事务可以统一提交,避免每次循环创建新连接导致的事务丢失。
- Upsert规范:更新时排除主键字段,使用
insert_stmt.excluded[col]引用本次插入的字段值,既符合SQL规范,也避免参数冗余。 - 空DataFrame判断:提前检查对比后的DataFrame是否为空,避免无意义的循环执行。
内容的提问来源于stack exchange,提问作者chibisuketyan
相关产品推荐
相关产品推荐

