You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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,但数据库始终无更新

问题原因

  1. 连接与事务管理错误:每次循环都创建新的engine.connect()上下文,上下文退出时连接自动关闭,事务未提交。最后一行的conn.commit()引用的是最后一次循环的连接,此时该连接已经被关闭,提交操作无效。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 04:40:17