如何提升Python将DataFrame导入SQL Server的速度?
优化Python向SQL Server批量插入速度的方案
针对你用SQLAlchemy上传13000行数据耗时3分钟的问题,以下是几个实用的提速方案,按改动成本从小到大排序:
1. 开启SQLAlchemy的fast_executemany参数
这是改动最小的方案,仅需在创建引擎时添加参数,就能大幅提升to_sql的插入速度。原理是让pyodbc使用ODBC的批量执行模式,减少网络交互次数。
修改后的代码:
import urllib.parse from sqlalchemy import create_engine quoted = urllib.parse.quote_plus("DRIVER={SQL Server};SERVER=......;DATABASE=.....") # 添加fast_executemany=True参数 engine = create_engine( 'mssql+pyodbc:///?odbc_connect={}'.format(quoted), fast_executemany=True ) try: print('uploading...') dataFrame.to_sql("MyTable", schema='dbo', con=engine, if_exists='replace', chunksize=1000) except Exception as e: print(f"上传失败: {e}") finally: engine.dispose()
2. 使用SQLAlchemy的bulk_insert_mappings
如果开启fast_executemany后仍达不到预期速度,可以尝试用SQLAlchemy的批量映射插入,减少ORM的额外开销。需要先定义对应的数据模型。
示例代码:
from sqlalchemy import create_engine, Column, String, Integer from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.orm import sessionmaker # 定义数据模型,需与MyTable结构一致 Base = declarative_base() class MyTableModel(Base): __tablename__ = 'MyTable' __table_args__ = {'schema': 'dbo'} # 替换为你的实际列名和类型 id = Column(Integer, primary_key=True) name = Column(String) # ...其他列 # 创建引擎和会话 quoted = urllib.parse.quote_plus("DRIVER={SQL Server};SERVER=......;DATABASE=.....") engine = create_engine('mssql+pyodbc:///?odbc_connect={}'.format(quoted), fast_executemany=True) Session = sessionmaker(bind=engine) session = Session() try: # 将DataFrame转为字典列表 data_mappings = dataFrame.to_dict('records') # 批量插入 session.bulk_insert_mappings(MyTableModel, data_mappings) session.commit() print('上传完成') except Exception as e: session.rollback() print(f"上传失败: {e}") finally: session.close() engine.dispose()
3. 直接使用pyodbc的executemany
绕开SQLAlchemy,直接用pyodbc底层操作,配合fast_executemany,速度更快。适合不需要ORM特性的场景。
示例代码:
import pyodbc import pandas as pd # 建立连接 conn = pyodbc.connect("DRIVER={SQL Server};SERVER=......;DATABASE=.....") cursor = conn.cursor() # 开启fast_executemany cursor.fast_executemany = True try: # 生成插入语句,自动匹配DataFrame的列 columns = ", ".join(dataFrame.columns) placeholders = ", ".join(["?"] * len(dataFrame.columns)) insert_sql = f"INSERT INTO dbo.MyTable ({columns}) VALUES ({placeholders})" # 将DataFrame转为列表的列表 data_list = dataFrame.values.tolist() # 执行批量插入 cursor.executemany(insert_sql, data_list) conn.commit() print('上传完成') except Exception as e: conn.rollback() print(f"上传失败: {e}") finally: cursor.close() conn.close()
4. 使用SQL Server原生BCP工具
如果数据量极大(比如百万级以上),SQL Server的BCP是最快的批量导入方式,通过Python调用命令行执行即可。需要确保系统已安装SQL Server命令行工具,且有足够权限。
示例代码:
import subprocess import tempfile import pandas as pd import os # 将DataFrame导出为临时CSV文件(用制表符分隔避免逗号冲突) with tempfile.NamedTemporaryFile(mode='w', delete=False, suffix='.csv', encoding='utf-8') as temp_file: dataFrame.to_csv(temp_file, index=False, sep='\t', na_rep='') # 构造BCP命令,替换为你的服务器、数据库、账号信息 bcp_command = [ 'bcp', 'dbo.MyTable in', temp_file.name, '-S', '你的服务器地址', '-d', '你的数据库名', '-U', '用户名', '-P', '密码', '-c', '-t\t', '-r\n', '-C', 'UTF-8' ] try: # 执行BCP命令 result = subprocess.run(bcp_command, check=True, capture_output=True, text=True) print(f"上传完成,BCP输出: {result.stdout}") except subprocess.CalledProcessError as e: print(f"BCP执行失败: {e.stderr}") finally: # 删除临时文件 os.unlink(temp_file.name)
内容的提问来源于stack exchange,提问作者Maxcot
相关产品推荐
相关产品推荐

