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

SQLAlchemy ORM实现SQL Server MP3记录UPSERT及mp3_id获取问题

问题解决:SQLAlchemy ORM同步MP3数据到SQL Server的存在性判断与主键提取问题

核心问题分析

  • 存在性判断逻辑错误:session.execute()返回的是Result对象,永远不会等于None,必须通过.scalar_one_or_none()等方法获取实际数据后再判断是否为空。
  • 主键提取失败:未正确获取ORM模型实例,直接操作Result对象无法访问mp3_id属性。

修正后的完整代码

from sqlalchemy import select, insert, update
from sqlalchemy.orm import sessionmaker
from sqlalchemy.exc import SQLAlchemyError
from datetime import datetime

# 假设engine、MP3_Files_indexed、brussels_tz、info_log、error_log已提前定义
i = 0
success_count = 0
failure_count = 0
Session = sessionmaker(bind=engine)
session = Session()

try:
    for sd in sql_data:
        # 通过path查询现有记录,直接返回ORM实例或None
        existing_record = session.execute(
            select(MP3_Files_indexed).where(MP3_Files_indexed.path == sd[1])
        ).scalar_one_or_none()

        if not existing_record:
            # 记录不存在,执行插入
            info_log.debug(f"mp3 file {sd[1]} not in DB\n")
            session.execute(
                insert(MP3_Files_indexed),
                [
                    {
                        "name":                sd[0],
                        "path":                sd[1],
                        "title":               sd[2],
                        "contributing_artist": sd[3],
                        "album":               sd[4],
                        "year":                sd[5],
                        "genre":               sd[6],
                        "track_num":           sd[7],
                        "number_of_tracks":    sd[8],
                        "disc_num":            sd[9],
                        "number_of_discs":     sd[10],
                        "created_on":          datetime.now(brussels_tz).isoformat(timespec='minutes'),
                        "updated_on":          datetime.now(brussels_tz).isoformat(timespec='minutes')
                    }
                ]
            )
        else:
            # 记录存在,用现有实例的mp3_id执行更新
            info_log.debug(f"mp3 file {sd[1]} is in DB\n")
            session.execute(
                update(MP3_Files_indexed)
                .where(MP3_Files_indexed.mp3_id == existing_record.mp3_id)
                .values(
                    name=sd[0],
                    path=sd[1],
                    title=sd[2],
                    contributing_artist=sd[3],
                    album=sd[4],
                    year=sd[5],
                    genre=sd[6],
                    track_num=sd[7],
                    number_of_tracks=sd[8],
                    disc_num=sd[9],
                    number_of_discs=sd[10],
                    updated_on=datetime.now(brussels_tz).isoformat(timespec='minutes')
                )
            )

        # 每处理100条记录提交一次事务,提升性能
        if i % 100 == 0:
            session.commit()
            print(f"Processed {i} records...")
        
        success_count += 1
        i += 1
    # 提交剩余未处理的记录
    session.commit()
    print(f"Total processed: {i} records. Success: {success_count}, Failure: {failure_count}")
except SQLAlchemyError as e:
    session.rollback()
    failure_count += 1
    error_log.error(f"Batch error: {str(e)}\n")
finally:
    session.close()

关键修改说明

  1. 存在性判断优化:

    • 使用scalar_one_or_none()方法,直接返回匹配的ORM模型实例(记录存在时)或None(记录不存在时),避免手动处理Result对象的冗余操作。
    • 用if not existing_record:替代原错误的if result == None:判断逻辑。
  2. 主键提取与更新逻辑修正:

    • 从existing_record直接访问mp3_id属性,因为它是完整的ORM模型实例。
    • 更新语句改用where()指定主键条件配合values()传递字段值,符合SQLAlchemy ORM的标准写法。
  3. 性能与稳定性优化:

    • 每100条记录提交一次事务,避免21800次单条提交带来的性能损耗。
    • 添加try-except-finally块,确保异常时回滚事务并关闭会话,防止数据库资源泄漏。
  4. 循环变量简化:

    • 原代码同时使用for sd in sql_data和i += 1访问sql_data[i][1],直接改用sd[1]更简洁,也避免了索引越界风险。

内容的提问来源于stack exchange,提问作者Jan J. Holvoet

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 04:25:57