使用Python将DataFrame写入MS SQL Server过慢,求优化方案
背景
我正在开发一个基于Python的平台,供非专业用户上传Excel数据到MS SQL Server数据库。用户选择Excel文件后,Python会解析生成多个DataFrame,再分别存入数据库对应的表中。
现状
当前处理一个含5万行、150列(总大小16MB)的Excel文件,生成12个DataFrame并入库需要2-3分钟;测试50MB文件时耗时长达7分钟,功能正常但速度远达不到预期。
已尝试方案
基础连接与DataFrame加载代码
# 连接字符串 connection_string = f""" DRIVER={{{DRIVER_NAME}}}; SERVER={{{SERVER_NAME}}}; DATABASE={{{DATABASE_NAME}}}; uid=XYZ; pwd=XYZ; Trust_Connection=yes; ColumnEncryption=Enabled; """ # 数据库连接 params=urllib.parse.quote_plus(connection_string) engine = sa.create_engine("mssql+pyodbc:///?odbc_connect={}".format(params), fast_executemany=True) con=engine.connect() # 读取Excel生成DataFrame df_Addr = pd.read_excel(excel_file, sheet_name = "Address_Details") df_Bank = pd.read_excel(excel_file, sheet_name = "Bank_Details") ... df_N = pd.read_excel(excel_file, sheet_name = "N_Details")
方案1:SQLAlchemy to_sql
# 存储地址表 saving_query_Address='DQ_Raw_Address' df_Addr.to_sql(saving_query_Address,engine,schema="dbo",if_exists='append',index=False, chunksize = 5000, dtype={'NAME1': sa.types.NVARCHAR(length=100), 'CITY1': sa.types.NVARCHAR(length=100), 'STREET': sa.types.NVARCHAR(length=100)}) # 存储银行表 saving_query_Bank='DQ_Raw_Bank' df_Bank.to_sql(saving_query_Bank,engine,schema="dbo",if_exists='append',index=False, chunksize = 5000, dtype={'_COMMENT':sa.types.VARCHAR(length=100),'_ACTION_CODE':sa.types.VARCHAR(length=100),'SOURCE_ID':sa.types.VARCHAR(length=100),'BKVID':sa.types.VARCHAR(length=100),'PARTNER':sa.types.VARCHAR(length=100),'BANKS':sa.types.VARCHAR(length=100),'IBAN':sa.types.VARCHAR(length=100),'ACCOUNT_ID':sa.types.VARCHAR(length=50),'CHECK_DIGIT':sa.types.VARCHAR(length=50),'ACCOUNT_TYPE':sa.types.VARCHAR(length=50),'BP_EEW_BUT0BK':sa.types.VARCHAR(length=50)}) # 其余10个表逻辑相同 # 总耗时:130秒
方案2:PyODBC executemany
# 存储地址表 saving_query_Address='DQ_Raw_Address' insert_to_tbl = f"INSERT INTO {saving_query_Address} VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)" cursor = conn.cursor() cursor.fast_executemany = True cursor.executemany(insert_to_tbl, df_Addr.values.tolist()) cursor.commit() cursor.close() # 存储银行表 saving_query_Bank='DQ_Raw_Bank' insert_to_tmp_tbl_stmt = f"INSERT INTO {saving_query_Bank} VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)" cursor = conn.cursor() cursor.fast_executemany = True cursor.executemany(insert_to_tmp_tbl_stmt, df_Bank.values.tolist()) cursor.commit() cursor.close() # 其余10个表逻辑相同 # 总耗时:200秒
备注:已尝试将Excel转CSV后加载DataFrame,无明显提速;无SQL Server Bulk Admin权限,无法使用BULK INSERT;需通过VPN连接数据库服务器。
版本信息:Pandas 1.5.0、PyODBC 4.0.34、SQLAlchemy 1.4.42
优化建议
1. 调优chunksize参数
当前chunksize=5000并非最优值,建议测试更大的块大小(如20000、50000),减少数据库交互次数。需根据服务器带宽和本地内存调整,避免块过大导致内存溢出。
2. 复用数据库连接与游标
方案2中每次插入都新建游标并关闭连接,会增加连接开销。建议复用同一连接和游标,完成所有插入后再关闭:
cursor = conn.cursor() cursor.fast_executemany = True # 依次插入所有表 insert_to_tbl = f"INSERT INTO DQ_Raw_Address VALUES (...)" cursor.executemany(insert_to_tbl, df_Addr.values.tolist()) insert_to_tmp_tbl_stmt = f"INSERT INTO DQ_Raw_Bank VALUES (...)" cursor.executemany(insert_to_tmp_tbl_stmt, df_Bank.values.tolist()) # 其余表插入逻辑... cursor.commit() cursor.close()
3. 批量提交事务
SQLAlchemy默认可能开启自动提交,建议手动控制事务,一次性提交所有插入操作,减少事务提交的网络往返开销:
with engine.begin() as conn: df_Addr.to_sql('DQ_Raw_Address', conn, schema="dbo", if_exists='append', index=False, chunksize=20000) df_Bank.to_sql('DQ_Raw_Bank', conn, schema="dbo", if_exists='append', index=False, chunksize=20000) # 其余表插入逻辑...
engine.begin()会自动管理事务,所有操作完成后一次性提交。
4. 对齐DataFrame与数据库字段类型
确保DataFrame列的数据类型与数据库表字段严格匹配,减少类型转换开销。例如限制字符串列长度与数据库字段一致:
df_Addr['NAME1'] = df_Addr['NAME1'].astype('string').str.slice(0, 100)
5. 并行插入(谨慎测试)
若服务器资源允许,可使用多线程并行处理不同DataFrame的插入,注意SQLAlchemy连接池配置(默认线程不安全,需设置pool_size和max_overflow):
from concurrent.futures import ThreadPoolExecutor def insert_df(df, table_name): with engine.connect() as conn: df.to_sql(table_name, conn, schema="dbo", if_exists='append', index=False, chunksize=20000) dfs_tables = [(df_Addr, 'DQ_Raw_Address'), (df_Bank, 'DQ_Raw_Bank'), ...] with ThreadPoolExecutor(max_workers=4) as executor: for df, tbl in dfs_tables: executor.submit(insert_df, df, tbl)
注意:VPN环境下并行可能增加网络负载,需测试验证是否提速。
6. 临时禁用索引与约束
若目标表存在非必要的索引或约束,可在插入前禁用,插入完成后重建,减少插入时的索引维护开销:
-- 禁用索引 ALTER INDEX ALL ON DQ_Raw_Address DISABLE; -- 插入数据 -- 重建索引 ALTER INDEX ALL ON DQ_Raw_Address REBUILD;
需注意权限与数据一致性,适合纯批量导入场景。
内容的提问来源于stack exchange,提问作者Arsalan Chaudhary

