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()
性能问题根源
- 单条数据频繁提交:遍历汇率数据时,每添加一条记录就调用一次
upsert_sql,每次都执行MERGE并提交事务,大量网络往返和事务开销直接拖垮性能。 - 未批量处理数据:当前代码每次只处理单个日期的单条汇率,没有积累到一定量再批量提交,导致SQL语句反复解析执行。
- 数据列表未重置:
upsert_list持续累积数据,每次调用upsert_sql都传入整个列表,导致SQL语句的VALUES部分越来越长,解析成本飙升。 - 未过滤已有日期:没有先查询数据库中已存在的日期,可能重复处理已有数据。
优化方案
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
相关产品推荐
相关产品推荐

