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

使用SQLAlchemy+pyodbc向SQL Server批量导入大CSV时遇内存错误

问题

使用Python脚本结合SQLAlchemy和pyodbc,将最大14GB的CSV文件批量导入本地SQL Server数据库。由于无法在内存中存储完整14GB的DataFrame,采用pandas分块功能进行批量插入,甚至尝试了仅100行的小批量,但加载过程中仍在同一节点触发内存错误。已移除目标表索引以降低负载,问题依旧。后续排查发现仅特定文件会触发该错误,单独加载该文件也会报错,怀疑错误提示可能不准确。

代码示例
import os
import glob
import traceback

import pandas as pd
import pyodbc
from sqlalchemy import create_engine, text
from sqlalchemy.exc import SQLAlchemyError
from tqdm import tqdm

# Replace the following variables with your database connection details
DB_USERNAME = '-----'
DB_PASSWORD = '-----'
DB_HOST = '------'
DB_NAME = '------'

DATA_DIRECTORY = "---------"
ERRORS_DIRECTORY = os.path.join(DATA_DIRECTORY, "errors")


def create_database_engine():
    engine = create_engine(f"mssql+pyodbc://{DB_USERNAME}:{DB_PASSWORD}@{DB_HOST}/{DB_NAME}?driver=ODBC+Driver+17+for+SQL+Server", fast_executemany=True )
    return engine

def batch_insert_to_database(engine, data, table, error_log_file):
    # Insert the data into the database in batches
    # The first batch will DROP ANY EXISTING TABLE IN THE DATABASE WITH THE SAME NAME
    try:
        data.to_sql(table, con=engine, index=False, if_exists='append')
    except SQLAlchemyError as e:
        error_message = f"Error during batch insertion: {e}"
        print("error encountered")
        
        with open(error_log_file, 'a') as error_log:
            error_log.write(error_message + '\n')
            traceback.print_exc(file=error_log)

        return False
    return True

def load_table_data(csv, table, errors_directory):
    
    print(f"Beginning load for {table}")
    
    # Create a database engine
    engine = create_database_engine()

    # Create an empty errors DataFrame to store failed batches
    errors_df = pd.DataFrame()

    # Batch size for insertion
    batch_size = 100

    # Initialize tqdm for progress tracking
    progress_bar = tqdm(total=0, desc="Processing")
    
    with engine.connect() as connection:
        truncate_command = text(f"TRUNCATE TABLE [dbo].[{table}]")
        connection.execute(truncate_command)

    error_report_file = os.path.join(errors_directory, f"errors_{table}.txt")
    # Read data from the CSV file in batches
    for batch_data in pd.read_csv(csv, chunksize=batch_size):
        # Try to insert the batch into the database
        success = batch_insert_to_database(engine, batch_data, table, error_report_file)

        # If the batch insertion fails, add the batch to the errors DataFrame
        if not success:
            errors_df = pd.concat([errors_df, batch_data])

        # Update the progress bar
        progress_bar.update(len(batch_data))

    # Close the progress bar
    progress_bar.close()


    error_data_file = os.path.join(errors_directory, f"errors_{table}.csv")
    # Save the errors DataFrame to a CSV file
    if not errors_df.empty:
        errors_df.to_csv(error_data_file, index=False)
        print(f"Errors saved to {error_data_file}")


def main():
    pattern = os.path.join(DATA_DIRECTORY, '**', '*.gz')
    gz_files = glob.glob(pattern, recursive=True)
    tables = [[file, file.split(os.path.sep)[-1].split(".")[0]] for file in gz_files]
    
    for table_data in tables:
        print(table_data)
        load_table_data(table_data[0], table_data[1], ERRORS_DIRECTORY)

if __name__ == "__main__":
    main()
报错栈信息
Traceback (most recent call last):
  File "C:\...\main.py", line 97, in <module>
    main()
  File "C:\...\main.py", line 94, in main
    load_table_data(table_data[0], table_data[1], ERRORS_DIRECTORY)
  File "C:\...\main.py", line 66, in load_table_data
    success = batch_insert_to_database(engine, batch_data, table, error_report_file)
              ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\...\main.py", line 29, in batch_insert_to_database
    data.to_sql(table, con=engine, index=False, if_exists='append')
  File "C:\...\pandas\core\generic.py", line 3008, in to_sql
    return sql.to_sql(
           ^^^^^^^^^^^
  File "C:\...\pandas\io\sql.py", line 788, in to_sql
    return pandas_sql.to_sql(
           ^^^^^^^^^^^^^^^^^^
  File "C:\...\pandas\io\sql.py", line 1958, in to_sql
    total_inserted = sql_engine.insert_records(
                     ^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\...\pandas\io\sql.py", line 1498, in insert_records
    return table.insert(chunksize=chunksize, method=method)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\...\pandas\io\sql.py", line 1059, in insert
    num_inserted = exec_insert(conn, keys, chunk_iter)
                   ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\...\pandas\io\sql.py", line 951, in _execute_insert
    result = conn.execute(self.table.insert(), data)
             ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\...\sqlalchemy\engine\base.py", line 1412, in execute
    return meth(
           ^^^^^
  File "C:\...\sqlalchemy\sql\elements.py", line 516, in _execute_on_connection
    return connection._execute_clauseelement(
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\...\sqlalchemy\engine\base.py", line 1635, in _execute_clauseelement
    ret = self._execute_context(
          ^^^^^^^^^^^^^^^^^^^^^^
  File "C:\...\sqlalchemy\engine\base.py", line 1844, in _execute_context
    return self._exec_single_context(
           ^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\...\sqlalchemy\engine\base.py", line 1984, in _exec_single_context
    self._handle_dbapi_exception(
  File "C:\...\sqlalchemy\engine\base.py", line 2342, in _handle_dbapi_exception
    raise exc_info[1].with_traceback(exc_info[2])
  File "C:\...\sqlalchemy\engine\base.py", line 1934, in _exec_single_context
    self.dialect.do_executemany(
  File "C:\...\sqlalchemy\dialects\mssql\pyodbc.py", line 716, in do_executemany
    super().do_executemany(cursor, statement, parameters, context=context)
  File "C:\...\sqlalchemy\engine\default.py", line 918, in do_executemany
    cursor.executemany(statement, parameters)

错误信息仅为MemoryError。

解决方案建议

1. 排查特定CSV文件的异常数据

既然仅特定文件触发错误,优先排查该文件内容:

  • 用pd.read_csv(csv, nrows=1000)读取前若干行,检查字段是否有超长文本、特殊字符或格式错乱;
  • 逐行读取文件,定位触发内存错误的具体行:
    import csv
    with open(csv, 'rt', encoding='utf-8') as f:
        reader = csv.reader(f)
        for idx, row in enumerate(reader):
            try:
                pd.DataFrame([row])
            except MemoryError:
                print(f"Error at row {idx}: {row}")
                break
    
  • 尝试手动解压.gz文件后再导入,排查是否存在压缩块损坏。

2. 优化内存管理与垃圾回收

代码中存在内存累积风险,做以下修改:

  • 用列表存储错误行,避免pd.concat反复生成新DataFrame:
    # 替换原errors_df初始化
    error_rows = []
    # 替换失败时的处理逻辑
    if not success:
        error_rows.append(batch_data)
    # 最后合并错误数据
    if error_rows:
        errors_df = pd.concat(error_rows, ignore_index=True)
    
  • 每个批次处理后强制释放内存:
    import gc
    # 在循环末尾添加
    del batch_data
    gc.collect()
    

3. 调整SQLAlchemy连接与插入参数

  • 复用同一个数据库连接,减少连接开销与资源占用:
    修改load_table_data函数,在循环外创建并保持连接:
    def load_table_data(csv, table, errors_directory):
        print(f"Beginning load for {table}")
        engine = create_database_engine()
        error_rows = []
        batch_size = 100
        progress_bar = tqdm(total=0, desc="Processing")
        
        # 复用同一个连接
        with engine.connect() as connection:
            # 执行truncate并提交
            truncate_command = text(f"TRUNCATE TABLE [dbo].[{table}]")
            connection.execute(truncate_command)
            connection.commit()
            
            error_report_file = os.path.join(errors_directory, f"errors_{table}.txt")
            for batch_data in pd.read_csv(csv, chunksize=batch_size):
                try:
                    batch_data.to_sql(table, con=connection, index=False, if_exists='append')
                    connection.commit()
                    success = True
                except SQLAlchemyError as e:
                    error_message = f"Error during batch insertion: {e}"
                    print("error encountered")
                    with open(error_report_file, 'a') as error_log:
                        error_log.write(error_message + '\n')
                        traceback.print_exc(file=error_log)
                    success = False
                
                if not success:
                    error_rows.append(batch_data)
                progress_bar.update(len(batch_data))
                
                del batch_data
                gc.collect()
        
        progress_bar.close()
        # 保存错误数据
        if error_rows:
            errors_df = pd.concat(error_rows, ignore_index=True)
            error_data_file = os.path.join(errors_directory, f"errors_{table}.csv")
            errors_df.to_csv(error_data_file, index=False)
            print(f"Errors saved to {error_data_file}")
    
  • 尝试关闭fast_executemany=True,该参数虽能加速插入,但部分复杂数据类型场景下会导致内存占用过高。

4. 优化CSV读取参数

  • 显式指定字段数据类型,避免pandas自动推断占用额外内存:
    # 示例:根据实际字段定义类型
    dtype_dict = {
        'id': 'int32',
        'text_field': 'string',
        'numeric_field': 'float32'
    }
    for batch_data in pd.read_csv(csv, chunksize=batch_size, dtype=dtype_dict):
        # 后续处理逻辑
    
  • 对于包含大文本的文件,添加low_memory=False参数,避免分块解析时的类型推断冲突。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 16:12:03