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

PYODBC实现SQL Upsert性能低下及相关技术问题求助

问题与代码优化方案

问题概况

现有Python代码功能正常但性能极差,每秒最多插入50行,后续甚至跌至不足10行。单独执行MERGE操作速度很快,但整体流程因额外步骤拖慢性能。核心需求:

  • 找出SQL_TABLE中2023年1月1日至今缺失的日期
  • 调用API获取对应日期的汇率数据并MERGE到表中
  • 获取最新汇率数据并MERGE到表中
  • 解决pyodbc中executemany()(含fastexecute)无法与MSSQL MERGE配合,以及MERGE语句无法完全参数化(只能用format()生成VALUES部分)的问题

原始代码

import datetime as dt
import requests
import pyodbc

# 假设已初始化数据库连接和游标
# cnxn = pyodbc.connect(...)
# crsr = cnxn.cursor()

def test_single_date():
    date = [dt.date(2023, 1, 1)]
    yield date

def get_data(url):
    api_response = requests.get(url)
    api_data = api_response.json()
    return api_data

def process_dates():
    # Step 1: 获取日期列表并创建SQL插入数据列表
    dates_to_process = test_single_date()
    upsert_list = []

    # Step 2: 生成URL并获取数据
    for date_list in dates_to_process:
        for single_date in date_list:

            # Step 2a: 生成URL
            base_hist_url = 'https://api_url/history/currency/{year}/{month}/{day}'
            hist_url = base_hist_url.format(year=single_date.year, month=single_date.month, day=single_date.day)

            # Step 2b: 获取API数据
            hist_api_data = get_data(hist_url)

            # Step 3: 批量准备数据并执行Upsert
            for k, v in hist_api_data['conversion_rates'].items():
                upsert_list.append(['currency_code', k, v, single_date])

                # Step 3a: 执行SQL Upsert
                upsert_sql(upsert_list)
    
    # Step 4: 获取最新汇率数据
    latest_url = 'https://api_url/latest/currency'
    lastest_api_data = get_data(latest_url)
    for k, v in lastest_api_data['conversion_rates'].items():
        upsert_list.append(['currency_code', k, v, dt.datetime.now(dt.UTC)])

        # Step 4a: 执行SQL语句
        upsert_sql(upsert_list)

def upsert_sql(upsert_list):
    sql = """
    MERGE INTO SQL_TABLE AS Target
    USING (
        VALUES {}
    ) AS Source (currency_from, currency_to, factor, timestamp)
    ON Target.currency_from = Source.currency_from 
    AND Target.currency_to = Source.currency_to
    AND CAST(Target.timestamp as date) = CAST(Source.timestamp as date)
    WHEN NOT MATCHED THEN
        INSERT (currency_from, currency_to, factor, timestamp) VALUES (Source.currency_from, Source.currency_to, Source.factor, Source.timestamp);
    """.format(','.join(['(?,?,?,?)' for _ in range(len(upsert_list))]))
    params = [item for sublist in upsert_list for item in sublist]

    try:
        crsr.execute(sql, params)
        crsr.commit()
    except Exception as e:
        crsr.rollback()
        print(e)
        print('Transaction rollback')

process_dates()
crsr.close()
cnxn.close()

性能问题根源

  1. 单条数据频繁提交:遍历汇率数据时,每添加一条记录就调用一次upsert_sql,每次都执行MERGE并提交事务,大量网络往返和事务开销直接拖垮性能。
  2. 未批量处理数据:当前代码每次只处理单个日期的单条汇率,没有积累到一定量再批量提交,导致SQL语句反复解析执行。
  3. 数据列表未重置:upsert_list持续累积数据,每次调用upsert_sql都传入整个列表,导致SQL语句的VALUES部分越来越长,解析成本飙升。
  4. 未过滤已有日期:没有先查询数据库中已存在的日期,可能重复处理已有数据。

优化方案

1. 批量提交数据

将单条提交改为批量提交,设定合理的批量大小(如1000条),减少事务和网络交互次数:

def process_dates():
    # 获取数据库中缺失的日期
    missing_dates = get_missing_dates()
    upsert_list = []
    batch_size = 1000  # 根据实际情况调整

    # 处理历史日期数据
    for single_date in missing_dates:
        base_hist_url = 'https://api_url/history/currency/{year}/{month}/{day}'
        hist_url = base_hist_url.format(year=single_date.year, month=single_date.month, day=single_date.day)
        hist_api_data = get_data(hist_url)
        
        for k, v in hist_api_data['conversion_rates'].items():
            upsert_list.append(['currency_code', k, v, single_date])
            # 达到批量大小就执行提交
            if len(upsert_list) >= batch_size:
                upsert_sql(upsert_list)
                upsert_list = []  # 重置列表
    # 提交剩余未达批量的数据
    if upsert_list:
        upsert_sql(upsert_list)
        upsert_list = []

    # 处理最新汇率数据
    latest_url = 'https://api_url/latest/currency'
    lastest_api_data = get_data(latest_url)
    for k, v in lastest_api_data['conversion_rates'].items():
        upsert_list.append(['currency_code', k, v, dt.datetime.now(dt.UTC)])
        if len(upsert_list) >= batch_size:
            upsert_sql(upsert_list)
            upsert_list = []
    if upsert_list:
        upsert_sql(upsert_list)

2. 查询缺失日期

新增函数获取2023-01-01至今数据库中缺失的日期,避免重复处理已有数据:

def get_missing_dates():
    # 查询已存在的日期
    sql = """
    SELECT DISTINCT CAST(timestamp AS DATE) AS existing_date
    FROM SQL_TABLE
    WHERE CAST(timestamp AS DATE) >= '2023-01-01'
    """
    crsr.execute(sql)
    existing_dates = {row[0] for row in crsr.fetchall()}

    # 生成2023-01-01至今的所有日期
    start_date = dt.date(2023, 1, 1)
    end_date = dt.date.today()
    all_dates = []
    current_date = start_date
    while current_date <= end_date:
        all_dates.append(current_date)
        current_date += dt.timedelta(days=1)
    
    # 返回缺失的日期
    return [date for date in all_dates if date not in existing_dates]

3. 优化MERGE语句参数化

pyodbc的executemany()确实无法直接适配MERGE语句,因为MERGE是单条语句处理多行数据,而executemany()是针对单条DML语句多次执行。当前用format()生成多行占位符的方式是安全的——占位符仅作为参数位置标记,实际参数通过execute()传入,不存在SQL注入风险。优化后的upsert_sql函数:

def upsert_sql(upsert_list):
    if not upsert_list:
        return
    # 生成对应数量的参数占位符
    placeholders = ','.join(['(?,?,?,?)' for _ in upsert_list])
    sql = f"""
    MERGE INTO SQL_TABLE AS Target
    USING (
        VALUES {placeholders}
    ) AS Source (currency_from, currency_to, factor, timestamp)
    ON Target.currency_from = Source.currency_from 
    AND Target.currency_to = Source.currency_to
    AND CAST(Target.timestamp AS DATE) = CAST(Source.timestamp AS DATE)
    WHEN NOT MATCHED THEN
        INSERT (currency_from, currency_to, factor, timestamp) 
        VALUES (Source.currency_from, Source.currency_to, Source.factor, Source.timestamp);
    """
    params = [item for sublist in upsert_list for item in sublist]
    try:
        crsr.execute(sql, params)
        crsr.commit()
    except Exception as e:
        crsr.rollback()
        print(f"错误信息: {e}")
        print('事务已回滚')

4. 其他辅助优化

  • 复用API连接:使用requests.Session()复用HTTP连接,减少API请求的握手开销。
  • 数据库索引优化:给SQL_TABLE的currency_from、currency_to、timestamp字段添加联合索引,加速MERGE的ON条件匹配速度。
  • 并发API请求:如果日期数量多,可使用aiohttp并发请求多个日期的API,减少等待时间。

关于executemany与MERGE的说明

MSSQL的MERGE语句设计为单条语句处理多行数据,而pyodbc的executemany()是针对单条INSERT/UPDATE语句重复执行,两者无法直接兼容。因此,用format()生成多行参数占位符,再通过execute()一次性传入所有参数的方式,是当前场景下的最优解,且完全安全。

内容的提问来源于stack exchange,提问作者LT_AKR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 18:27:05