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

增量抽取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-profiler
    
    在目标函数添加装饰器:
    from memory_profiler import profile
    
    @profile
    def main():
        # ais_eod_fe的main函数逻辑
    

3. 优化程序运行逻辑

  • 减少重复操作:
    cifextraction_eod和accextraction_eod若为静态数据,可只抽取一次,无需每次启动都重新执行。
  • 细化日志监控:
    在bftextraction_initial中单独统计数据库查询耗时、返回行数,精准定位瓶颈环节。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:57:04