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

如何通过Python从Mac高效向远程SQL Server上传超大规模CSV数据

高效从Python向远程SQL Server批量导入CSV数据的优化方案

问题背景

需要从Mac通过Python向远程Microsoft SQL Server上传1000个CSV文件:

  • 每个文件含5列、1000万行数据,其中3列为float(长度≤8位),2列为string(长度≤30字符)
  • 当前SQLAlchemy方案效率极低:
    • .to_sql()分块插入需将chunksize设为400,整体耗时长达数月
    • SQLAlchemy execute批量插入单文件仍需20分钟,总耗时约333小时
  • 限制:仅使用Python实现,排除其他语言/服务器内置工具,已确认网速、服务器负载无问题

优化方案

1. 启用pyodbc的fast_executemany参数(SQL Server专属核心优化)

SQL Server的pyodbc驱动提供了fast_executemany参数,开启后会采用高效的批量参数绑定模式,大幅减少网络交互次数,可将插入速度提升数倍。配合SQLAlchemy使用:

from sqlalchemy import create_engine
import pandas as pd

# 构造连接字符串,替换为你的服务器信息
connection_string = (
    "DRIVER={ODBC Driver 18 for SQL Server};"
    "SERVER=your_server_address;"
    "DATABASE=your_db_name;"
    "UID=your_username;"
    "PWD=your_password;"
    "Encrypt=yes;"
    "TrustServerCertificate=yes;"
)

# 创建引擎时开启fast_executemany
engine = create_engine(
    "mssql+pyodbc:///?odbc_connect=" + connection_string,
    connect_args={"fast_executemany": True}
)

def insert_csv_with_sqlalchemy(csv_path):
    # 指定字段类型,避免pandas自动推断带来的性能损耗
    dtype_map = {
        'INSTRUMENT': 'string',
        'BID': 'float64',
        'ASK': 'float64',
        'MID_PRICE': 'float64',
        'TIMESTAMP_UTC_YYYYMMDDhhmmss_sss': 'string'
    }
    # 分块读取CSV,降低内存占用
    for chunk in pd.read_csv(csv_path, dtype=dtype_map, chunksize=100000):
        chunk.to_sql(
            'quotes_test',
            con=engine,
            if_exists='append',
            index=False,
            method='multi'
        )
        print(f"完成 {csv_path} 中 {len(chunk)} 行数据上传")

2. 直接使用pyodbc批量插入(绕过SQLAlchemy中间层)

若SQLAlchemy的封装仍有性能损耗,直接用pyodbc操作可进一步降低开销,核心仍是开启fast_executemany:

import pyodbc
import pandas as pd

def pyodbc_bulk_insert(csv_path):
    # 替换为你的服务器连接信息
    conn = pyodbc.connect(
        "DRIVER={ODBC Driver 18 for SQL Server};"
        "SERVER=your_server_address;"
        "DATABASE=your_db_name;"
        "UID=your_username;"
        "PWD=your_password;"
        "Encrypt=yes;"
        "TrustServerCertificate=yes;"
    )
    cursor = conn.cursor()
    cursor.fast_executemany = True  # 开启高效批量模式

    insert_sql = """
        INSERT INTO quotes_test (INSTRUMENT, BID, ASK, MID_PRICE, TIMESTAMP_UTC_YYYYMMDDhhmmss_sss)
        VALUES (?, ?, ?, ?, ?)
    """

    dtype_map = {
        'INSTRUMENT': 'string',
        'BID': 'float64',
        'ASK': 'float64',
        'MID_PRICE': 'float64',
        'TIMESTAMP_UTC_YYYYMMDDhhmmss_sss': 'string'
    }

    # 分块读取并插入
    for chunk in pd.read_csv(csv_path, dtype=dtype_map, chunksize=100000):
        # 将DataFrame转换为元组列表,适配pyodbc参数格式
        params = [tuple(row) for row in chunk.itertuples(index=False)]
        cursor.executemany(insert_sql, params)
        conn.commit()
        print(f"完成 {csv_path} 中 {len(chunk)} 行数据上传")

    cursor.close()
    conn.close()

3. 多进程并行处理多个CSV文件

利用Mac的多核CPU,同时处理多个CSV文件,可大幅缩短总耗时:

import os
from multiprocessing import Pool

def process_single_csv(csv_path):
    # 调用上述任意一个插入函数
    pyodbc_bulk_insert(csv_path)

if __name__ == "__main__":
    csv_directory = "/path/to/your/csv/folder"
    # 获取所有CSV文件路径
    csv_files = [
        os.path.join(csv_directory, filename) 
        for filename in os.listdir(csv_directory) 
        if filename.endswith(".csv")
    ]
    # 根据CPU核心数设置进程数(建议设为核心数的1-2倍)
    with Pool(processes=4) as pool:
        pool.map(process_single_csv, csv_files)

额外优化建议

  • 预匹配表结构:确保目标表quotes_test的字段类型与CSV完全一致,例如string类型设为VARCHAR(30),float类型设为FLOAT,避免插入时的类型转换开销。
  • 优化事务提交:每插入10万-100万行后提交一次事务,减少事务日志的写入压力。
  • 精简CSV读取:若不需要DataFrame的功能,可直接使用Python内置csv模块读取文件,进一步降低内存占用和处理时间。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 03:14:56