SQLAlchemy无法正确更新列,添加session.merge()触发主键冲突
问题
我编写了一段SQLAlchemy脚本,但未达到预期效果:
for index, row in df_journey.iterrows(): journey = JourneySummary(**row) print(vars(journey)) existing_journey = session.query(JourneySummary).filter( JourneySummary.RouteIDs == journey.RouteIDs, JourneySummary.Planned_Start_Time == journey.Planned_Start_Time ).first() if existing_journey is None: session.add(journey) added_rows += 1 else: # Update the existing journey with the new data for attr, value in vars(journey).items(): setattr(existing_journey, attr, value) merged_rows += 1 try: session.commit() print("Data successfully updated") except SQLAlchemyError as e: session.rollback() print(f"Data update failed: {e}") print(str(e))
问题在于当RouteIDs和Planned_Start_Time匹配条件时,脚本并未更新这些已存在行的列数据。
我尝试在如下位置添加session.merge():
else: # Update the existing journey with the new data for attr, value in vars(journey).items(): setattr(existing_journey, attr, value) session.merge() merged_rows += 1
但出现以下错误:
IntegrityError: ('23000', "[23000] [Microsoft][ODBC Driver 18 for SQL Server][SQL Server]Violation of PRIMARY KEY constraint 'PK__Journeys__3214EC07F297BA5F'. Cannot insert duplicate key in object 'dbo.Journeys'. The duplicate key value is (2860). (2627) (SQLParamData); [23000] [Microsoft][ODBC Driver 18 for SQL Server][SQL Server]The statement has been terminated. (3621)")
这表明session.merge()函数即使在条件匹配时仍尝试创建新行,求解决建议。
解决建议
1. 移除无效的session.merge()调用
session.merge()必须传入对象参数才能生效,你直接调用session.merge()属于错误用法。而且你已经从数据库加载了existing_journey,它本身处于会话管理状态,修改属性后直接提交会话就会自动同步到数据库,完全不需要额外调用merge。
2. 修复属性遍历的问题
vars(journey)会包含SQLAlchemy对象的内部属性(比如_sa_instance_state),这些属性不应该被赋值给数据库模型。应该只遍历模型定义的字段:
# 获取模型的所有数据库列名 model_columns = [col.name for col in JourneySummary.__table__.columns] for attr, value in journey.__dict__.items(): if attr in model_columns: setattr(existing_journey, attr, value)
或者用更高效的批量更新方式:
else: # 直接通过查询语句批量更新,避免遍历属性 session.query(JourneySummary).filter( JourneySummary.RouteIDs == journey.RouteIDs, JourneySummary.Planned_Start_Time == journey.Planned_Start_Time ).update(row.to_dict()) merged_rows += 1
注:如果row是pandas Series,to_dict()会自动转换为字段名-值的字典,刚好匹配模型字段。
3. 优化会话提交时机
不要在循环内每次迭代都提交会话,批量处理完成后统一提交,既能提升性能,也能减少数据库连接开销:
try: for index, row in df_journey.iterrows(): journey = JourneySummary(**row) existing_journey = session.query(JourneySummary).filter( JourneySummary.RouteIDs == journey.RouteIDs, JourneySummary.Planned_Start_Time == journey.Planned_Start_Time ).first() if existing_journey is None: session.add(journey) added_rows += 1 else: model_columns = [col.name for col in JourneySummary.__table__.columns] for attr, value in journey.__dict__.items(): if attr in model_columns: setattr(existing_journey, attr, value) merged_rows += 1 # 循环结束后统一提交 session.commit() print(f"数据更新完成:新增{added_rows}行,更新{merged_rows}行") except SQLAlchemyError as e: session.rollback() print(f"数据更新失败:{e}") print(str(e))
4. 正确理解session.merge()的适用场景
如果要使用merge,需要传入目标对象,它会通过主键判断对象是否存在:
# 替换原有的if-else逻辑 merged_journey = session.merge(journey) if merged_journey in session.new: added_rows += 1 else: merged_rows += 1
但这种方式依赖主键匹配,而你的业务唯一键是RouteIDs+Planned_Start_Time(非主键),所以会因为主键重复报错,这也是你之前出错的核心原因。因此针对你的业务场景,还是用查询+更新的逻辑更可靠。
内容的提问来源于stack exchange,提问作者Rexxology

