使用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
相关产品推荐
相关产品推荐

