Python写入MSSQL数据库性能优化求助:提升插入速度
问题描述
我有一段Python脚本,通过以下方式读取文件:
with open(filepath + filename, 'r') as inFile: for line in inFile: ...
该脚本读取并处理7GB的多列文件仅需约3分钟,但向MSSQL数据库插入数据的速度极慢。
我尝试过这些方法但效果不佳:
- 模仿SSIS数据流插入但失败(不清楚底层实现)
- pyodbc的常规
execute、executemany方法,即使启用fast_executemany速度仍不理想 - 使用bcpandas批量插入可达约5000行/秒,但将数据转换为DataFrame的耗时占总流程近一半
- 尝试通过打印数据让SSIS执行插入,但Python打印效率过低
- 生成新文件供SSIS读取的耗时与
executemany相当
请问还有哪些遗漏的优化方法?SSIS数据流采用何种方式实现数据库插入?
相关代码如下:
import pyodbc # tested executemany() import time import pandas as pd import bcpandas import sqlalchemy # start time startTime = time.time() ## file parameters filepath = "C:/" filename = "file.txt" alchemy_eng = sqlalchemy.create_engine('mssql+pyodbc:///?odbc_connect=DRIVER={ODBC Driver 17 for SQL Server};Server=localhost;Database=test;Trusted_Connection=yes;') bcpandas_eng = bcpandas.SqlCreds.from_engine(alchemy_eng) ### process file counter = 0 insert_params = [] with open(filepath + filename, 'r') as infile: for line in inFile: # temp table partial structure col1,col2,col3,.. = line.split('|') if col2 == "valueA": col2 = None # couple more .. col3 = col3.replace(',', '.') # couple more .. # commit every million not all to avoid memory crash if counter == 1000000: df = pd.DataFrame(insert_params) df.columns = ['col1',...] bcpandas.to_sql(df, 'table1', bcpandas_eng, if_exists="append") # reset insert_params = [] counter = 0 break if counter == 1: cursor.execute(insert_query, insert_params) break # list of elements insert_params.append((col1,col2,col3,..),) counter += 1 # end time endTime = time.time() # elapsed time print("elapsed time", endTime - startTime)
遗漏的优化方法
1. 直接调用SQL Server BCP工具(跳过DataFrame转换)
既然bcpandas的瓶颈在DataFrame转换,不如直接用SQL Server自带的bcp命令行工具。Python只负责处理数据并写入临时文件,再通过subprocess调用bcp批量导入,完全跳过DataFrame转换步骤:
import subprocess # 处理数据并写入临时文件 with open(filepath + filename, 'r') as infile, open('temp_bcp.txt', 'w') as temp_file: counter = 0 for line in infile: col1, col2, col3 = line.split('|') # 数据处理逻辑 if col2 == "valueA": col2 = '' # 适配BCP空值格式 col3 = col3.replace(',', '.') # 用制表符分隔字段(可根据表结构调整分隔符) temp_file.write(f"{col1}\t{col2}\t{col3}\n") # 每10万行刷新一次,避免内存占用过高 if counter % 100000 == 0: temp_file.flush() counter += 1 # 调用bcp命令导入 bcp_cmd = [ 'bcp', 'test.dbo.table1', 'in', 'temp_bcp.txt', '-S', 'localhost', '-T', '-c', '-t\t', '-r\n', '-k' # -k保留空值 ] subprocess.run(bcp_cmd, check=True)
2. 优化pyodbc的fast_executemany配置
确保正确启用fast_executemany,并调整批量大小(建议测试5000-20000行的批次),同时手动控制事务提交减少开销:
import pyodbc conn = pyodbc.connect('DRIVER={ODBC Driver 17 for SQL Server};Server=localhost;Database=test;Trusted_Connection=yes;') cursor = conn.cursor() cursor.fast_executemany = True batch_size = 15000 # 可根据实际环境调整 insert_query = "INSERT INTO table1 (col1, col2, col3) VALUES (?, ?, ?)" with open(filepath + filename, 'r') as infile: counter = 0 params = [] for line in infile: col1, col2, col3 = line.split('|') if col2 == "valueA": col2 = None col3 = col3.replace(',', '.') params.append((col1, col2, col3)) counter += 1 if counter % batch_size == 0: cursor.executemany(insert_query, params) conn.commit() params = [] # 提交剩余数据 if params: cursor.executemany(insert_query, params) conn.commit() conn.close()
3. 使用SQL Server BULK INSERT语句
如果处理后的文件格式符合要求,可直接用BULK INSERT从文件导入,效率远高于逐行插入:
from sqlalchemy import create_engine alchemy_eng = create_engine('mssql+pyodbc:///?odbc_connect=DRIVER={ODBC Driver 17 for SQL Server};Server=localhost;Database=test;Trusted_Connection=yes;') # 先生成符合格式的临时文件(同BCP思路) # 执行BULK INSERT with alchemy_eng.connect() as conn: conn.execute(""" BULK INSERT test.dbo.table1 FROM 'C:/temp_bcp.txt' WITH ( FIELDTERMINATOR = '\t', ROWTERMINATOR = '\n', KEEPNULLS, TABLOCK # 启用表锁提升速度 ) """) conn.commit()
4. 避免不必要的内存占用
当前代码中insert_params会缓存百万行数据,可适当减小批次大小(比如20万行),同时处理完一行就丢弃原始行数据,避免内存累积。
SSIS数据流的插入实现方式
SSIS数据流任务批量插入MSSQL主要依赖两种核心机制:
- OLE DB快速加载(IRowsetFastLoad):当OLE DB目标启用“快速加载”选项时,SSIS会调用OLE DB的
IRowsetFastLoad接口,直接将内存中的数据块写入数据库,同时会:- 禁用约束检查、触发器(可选)
- 使用批量日志模式(需数据库处于简单/大容量日志恢复模式)
- 减少网络往返和事务提交次数
- BCP底层逻辑集成:处理平面文件源时,SSIS会直接复用BCP的高效数据读取/写入管道,避免中间数据转换开销,以流的形式批量导入数据。
此外,SSIS通过缓冲区管理优化内存使用,将数据分成多个缓冲区并行处理,同时利用多线程执行数据转换和加载,进一步提升整体效率。
内容的提问来源于stack exchange,提问作者Sleepy
相关产品推荐
相关产品推荐

