增量抽取MSSQL交易数据耗时递增问题求助(附代码)
MSSQL按日增量抽取数据耗时随日期递增的排查与解决
问题背景
我在MSSQL数据库中有2022-01-01至2022-12-31的交易数据,通过pyodbc按日增量抽取,之后做预处理、特征工程并保存为CSV。但抽取耗时随日期递增:2022-01-01仅需10秒,2022-01-29却超过3小时。目前临时方案是每日21点停止Python程序,次日通过任务计划从checkpoint重启。
代码片段
main.py
def should_stop_execution(): current_time = datetime.now().time() stop_time = datetime.strptime('21:00', '%H:%M').time() return current_time >= stop_time # Checkpoint Functions def save_checkpoint(last_processed_date, checkpoint_file='last_processed_date.txt'): with open(checkpoint_file, 'w') as f: f.write(last_processed_date.strftime('%Y%m%d')) def load_checkpoint(checkpoint_file='last_processed_date.txt'): try: with open(checkpoint_file, 'r') as f: return datetime.strptime(f.read().strip(), '%Y%m%d') except FileNotFoundError: return None @log_to_file('AIS_log.log') def main(server, database, table, start_date_str, end_date_str): cif_path = "./Data/raw/cif" acc_path = "./Data/raw/acc" utils.delete_folder_contents(cif_path) ais_read_save_data.cifextraction_eod(server, database, cif_table) utils.delete_folder_contents(acc_path) ais_read_save_data.accextraction_eod(server, database, acc_table) # Attempt to load from checkpoint checkpoint_start_date = load_checkpoint() if checkpoint_start_date: start_date = checkpoint_start_date else: start_date = pd.to_datetime(start_date_str, format='%Y%m%d') end_date = pd.to_datetime(end_date_str, format='%Y%m%d') print(start_date) print(end_date) all_time = [] try: while start_date <= end_date: # Check if it's time to stop if should_stop_execution(): print("Stopping execution at", datetime.now()) break start = time.time() txn_path = "./data/raw/txn/" utils.delete_folder_contents(txn_path) txn = ais_read_save_data.bftextraction_initial(server, database, table, start_date) # Check if it's time to stop if should_stop_execution(): print("Stopping execution at", datetime.now()) break print('read save complete') if txn is not None: del txn # Start processing - Assuming ais_eod_fe.main() is another process you want to call ais_eod_fe.main() end = time.time() - start print(end) all_time.append((str(start_date), end)) # Save checkpoint after successful day processing save_checkpoint(start_date) # Move to next date start_date += timedelta(days=1) except Exception as e: print(e) finally: df = pd.Series(all_time) df.to_csv('timelog.csv', index=False) if __name__ == '__main__': main(server, database, bftranhist, bftranhist_startdate, bftranhist_enddate)
extraction.py
def bftextraction_initial(server, database, table, start_date, save_to_raw=True): parquet_path = f"{parent_dir}/data/raw/txn/" if not user: SQL_SERVER_ENGINE_URL = f"mssql+pyodbc:///?odbc_connect={urllib.parse.quote_plus('DRIVER={SQL Server};SERVER=' + server + ';DATABASE=' + database + ';Trusted_Connection=yes;')}" else: SQL_SERVER_ENGINE_URL = f"mssql+pyodbc:///?odbc_connect={urllib.parse.quote_plus(f'DRIVER={{SQL Server}};SERVER={server};DATABASE={database};UID={user};PWD={password};')}" engine = create_engine(SQL_SERVER_ENGINE_URL) columns_to_select = [ 'TRAN_NO', 'SEQ_AUTO', 'ACC_NO', 'NO_CIF', 'TRAN_PDATE', 'TRAN_CODE1', 'PROD_TYPE', 'TRAN_LOC', 'AMT_CR', 'AMT_DR', 'TRAN_TYPE', 'AMT_CUR', 'AMT_RATE', 'SEND_NAME', 'BENE_NAME', 'TRAN_CODE2', 'TRAN_CHAN', 'SEND_BCTY', 'BENE_BCTY', 'UNIT_PRICE' ] columns_str = ', '.join(columns_to_select) end_date_int = datetime.strftime(start_date, '%Y%m%d') print(end_date_int) # REMOVE TOP 100000 AFTER TEST stmt = text(f"SELECT {columns_str} from {urllib.parse.quote_plus(table)} where TRAN_PDATE = {end_date_int}") try: with engine.connect() as conn: df = pd.read_sql(stmt, conn) if df is not None: print(df.shape) df = df.drop_duplicates(['TRAN_NO', 'SEQ_AUTO', 'ACC_NO', 'NO_CIF', 'TRAN_PDATE']) # print(df.shape) if save_to_raw: dh.df_to_parquet(df, parquet_path) return df finally: engine.dispose() gc.collect()
可能原因分析
- 数据库查询性能瓶颈:
TRAN_PDATE字段无索引:按日期过滤时触发全表扫描,随着后续日期交易数据量增加,扫描耗时呈指数增长。- 隐式类型转换:
TRAN_PDATE若为日期类型,查询时用整数格式(如20220101)匹配,会触发隐式转换导致索引失效,强制全表扫描。
- Python内存泄漏:
ais_eod_fe.main()可能存在未释放的内存对象,导致程序运行越久内存占用越高,GC和IO操作变慢。 - 数据库连接重复创建:
每次抽取都新建数据库引擎,连接创建与销毁的开销累积,同时可能存在连接未彻底回收,导致数据库连接数饱和。
解决方案
1. 优化数据库查询
- 添加覆盖索引:
在MSSQL中执行以下语句,避免全表扫描和书签查找:CREATE NONCLUSTERED INDEX IX_BFTranHist_TRAN_PDATE ON bftranhist(TRAN_PDATE) INCLUDE (TRAN_NO, SEQ_AUTO, ACC_NO, NO_CIF, TRAN_CODE1, PROD_TYPE, TRAN_LOC, AMT_CR, AMT_DR, TRAN_TYPE, AMT_CUR, AMT_RATE, SEND_NAME, BENE_NAME, TRAN_CODE2, TRAN_CHAN, SEND_BCTY, BENE_BCTY, UNIT_PRICE); - 修正日期查询逻辑:
使用参数化查询匹配日期类型,避免隐式转换:# 替换extraction.py中的stmt部分 stmt = text(f"SELECT {columns_str} from {urllib.parse.quote_plus(table)} where TRAN_PDATE = :date_val") df = pd.read_sql(stmt, conn, params={"date_val": start_date.date()})
2. 优化Python资源管理
- 复用数据库引擎:
在main函数初始化一次引擎,传递给抽取函数,减少连接开销:# main.py开头添加引擎初始化 if not user: SQL_SERVER_ENGINE_URL = f"mssql+pyodbc:///?odbc_connect={urllib.parse.quote_plus('DRIVER={SQL Server};SERVER=' + server + ';DATABASE=' + database + ';Trusted_Connection=yes;')}" else: SQL_SERVER_ENGINE_URL = f"mssql+pyodbc:///?odbc_connect={urllib.parse.quote_plus(f'DRIVER={{SQL Server}};SERVER={server};DATABASE={database};UID={user};PWD={password};')}" engine = create_engine(SQL_SERVER_ENGINE_URL) # 调用抽取函数时传递engine txn = ais_read_save_data.bftextraction_initial(server, database, table, start_date, engine=engine) # 修改extraction.py函数定义 def bftextraction_initial(server, database, table, start_date, engine=None, save_to_raw=True): if engine is None: # 保留原引擎创建逻辑作为降级方案 ... # 移除finally中的engine.dispose(),改在main的finally中统一释放 - 排查内存泄漏:
使用memory-profiler监控ais_eod_fe.main()的内存使用:
在目标函数添加装饰器:pip install memory-profilerfrom memory_profiler import profile @profile def main(): # ais_eod_fe的main函数逻辑
3. 优化程序运行逻辑
- 减少重复操作:
cifextraction_eod和accextraction_eod若为静态数据,可只抽取一次,无需每次启动都重新执行。 - 细化日志监控:
在bftextraction_initial中单独统计数据库查询耗时、返回行数,精准定位瓶颈环节。
内容的提问来源于stack exchange,提问作者AIS TESS
相关产品推荐
相关产品推荐

