Python调用MySQL存储过程批量插入数据性能优化求助
问题:批量CSV数据插入MySQL性能瓶颈优化
原始问题
每月需多次读取客户提供的分隔符格式文件,行数从数百到百万不等。当前已实现通过Python逐行读取文件并调用MySQL存储过程sp_InsertMetaData插入数据,但性能极差:原接口引擎处理5000行耗时1小时,Python版本即使仅准备存储过程调用,处理2万行也需20分钟。
需求:寻求最优的多线程(或其他)优化方案,以提升文件读取及数据插入数据库的速度,重点优化存储过程调用环节。
相关代码:
def ReadMetaDataFile (srcFile, fileDelimiter, clientID, mysqlCnx): functionResponse = '' path, filename = os.path.split(srcFile) try: with open(srcFile, 'r') as delimitedFile: fileReader = csv.DictReader(delimitedFile, delimiter=fileDelimiter) #next(fileReader, None) if(mysqlCnx and mysqlCnx.is_connected()): cursor = mysqlCnx.cursor() for row in fileReader: rowArgs = (clientID, row['FILENAME'], '', row['DATA_FORMAT'], 1) cursor.callproc('sp_InsertMetaData', rowArgs) mysqlCnx.commit() cursor.close() mysqlCnx.close() functionResponse = "MetaData File Loaded Successfully" except mysql.connector.Error as err: functionResponse = "Error processing MySQL SPROC call: " + str(err) except (mysql.connector.Error, IOError) as err: functionResponse = "Error connecting to MySQL in function: " + str(err) except Exception as err: functionResponse = "Exception Found in function: " + str(err) finally: mysqlCnx.close() return functionResponse
更新内容
已尝试多种方式,目前仅构建存储过程调用就耗时过长,处理2万行需20分钟,百万级文件完全无法满足需求,恳请提供优化建议。
更新后代码:
def ReadMetaDataFile (srcFile, fileDelimiter, clientID, engine): functionResponse = '' path, filename = os.path.split(srcFile) try: query = 'CALL sp_InsertMetaData (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)' with open(srcFile, 'r') as delimitedFile: fileReader = csv.DictReader(delimitedFile, delimiter=fileDelimiter) conn = engine.raw_connection() cursor = conn.cursor() for row in fileReader: rowArgs = (clientID, row['FILENAME'], '', row['DATA_FORMAT'], row['SOURCE_SYSTEM_NAME'], row['CONFIDENTIALITY_CODE'], row['STATUS'], row['DOCUMENT_TYPE_SOURCE_SYSTEM'], row['DOCUMENT_TYPE_ID'], row['DOCUMENT_TYPE_DESCRIPTION'], row['DATE_OF_SERVICE'], row['DOCUMENT_ID'], row['SOURCE_CREATED_DATE'], row['SOURCE_LAST_MODIFIED_DATE'], row['TIMEZONE'], row['MRNSOURCE_SYSTEM'], row['PATIENT_MRN'], row['MEMBER_NBR'], row['PATIENT_LAST_NAME'], row['PATIENT_FIRST_NAME'], row['PATIENT_MIDDLE_NAME'], row['GENDER'], row['PATIENT_DATE_OF_BIRTH'], row['ENCOUNTER_SOURCE'], row['ENCOUNTER_ID'], row['ENCOUNTER_TYPE'], row['ADMIT_TIME'], row['DISCHARGE_TIME'], row['FACILITY_NAME'], row['FACILITY_SOURCE_SYSTEM'], row['PROVIDER_TYPE'], row['PROVIDER_SOURCE_SYSTEM'], row['PROVIDER_IDENTIFIER'], row['PROVIDER_LAST_NAME'], row['PROVIDER_FIRST_NAME'], row['PROVIDER_MIDDLE_NAME'], row['PROVIDER_CREDENTIAL'], row['PROVIDER_SPECIALTY'], 1) cursor.callproc("sp_InsertMetaData", rowArgs) #conn.commit() cursor.close() functionResponse = "MetaData File Loaded Successfully"
优化建议
1. 核心优化:替换逐行存储过程调用
存储过程单调用开销极大,逐行调用是性能瓶颈的根源。按以下方向调整:
- 如果存储过程仅做单表插入:直接改用批量INSERT语句,比如
INSERT INTO table (col1, col2...) VALUES (...), (...), (...),每次批量提交1000-5000行(根据单条数据大小调整)。 - 如果存储过程包含复杂业务逻辑:要么将逻辑迁移到Python批量处理后插入,要么修改存储过程支持批量参数输入(MySQL 8.0+支持表值参数,旧版本可通过分隔字符串拆分实现,但性能仍不如直接批量插入)。
2. 减少事务提交次数
原始代码每一行都执行commit,会触发大量磁盘IO,严重拖慢速度:
- 批量处理N行后统一调用一次
commit,比如每1000行提交一次。 - 加入异常捕获,批量提交失败时回滚当前批次并记录错误行,保证数据可靠性。
3. 提升CSV读取效率
- 用
csv.reader替代csv.DictReader:后者构建字典的开销远大于列表,百万级数据下差异明显,直接按列索引取数据即可。 - 打开文件时指定大缓冲区:比如
open(srcFile, 'r', buffering=1024*1024),减少磁盘IO次数。
4. 数据库连接与游标优化
- 使用连接池:复用数据库连接,避免频繁创建/关闭连接的开销(比如
mysql.connector的pooling模块或SQLAlchemy连接池)。 - 每个线程/进程持有独立连接:MySQL连接线程不安全,多线程环境下必须保证连接独占。
5. 多线程/多进程的正确用法
当前场景属于IO密集型(文件读取+网络请求),优先用线程池,超大文件可考虑多进程:
- 线程池:用
concurrent.futures.ThreadPoolExecutor,线程数设置为4-8(根据CPU核心数和数据库连接池上限调整)。 - 多进程:用
multiprocessing.Pool,将文件分割成多个块,每个进程处理一块,注意进程间数据隔离。
6. 数据库层面优化
- 临时关闭索引:插入大量数据前执行
ALTER TABLE table_name DISABLE KEYS,插入完成后执行ALTER TABLE table_name ENABLE KEYS,避免索引维护的额外开销。 - 关闭自动提交:执行
SET autocommit=0,手动批量提交。 - 调整MySQL配置:增大
innodb_buffer_pool_size、innodb_log_file_size等参数,提升写入性能。
示例代码(批量插入+线程池)
import csv import os from concurrent.futures import ThreadPoolExecutor import mysql.connector from mysql.connector import pooling # 初始化数据库连接池 db_pool = mysql.connector.pooling.MySQLConnectionPool( pool_name="mypool", pool_size=4, host="your_host", user="your_user", password="your_pass", database="your_db" ) def process_batch(batch_data): conn = db_pool.get_connection() cursor = conn.cursor() try: # 替换为实际表名和列顺序 insert_sql = """ INSERT INTO metadata_table (client_id, filename, col3, data_format, source_system_name, confidentiality_code, status, flag) VALUES (%s, %s, %s, %s, %s, %s, %s, %s) """ cursor.executemany(insert_sql, batch_data) conn.commit() except Exception as e: conn.rollback() raise e finally: cursor.close() conn.close() def process_file_chunk(chunk_data, clientID): batch_size = 1000 batch = [] for row in chunk_data: row_args = ( clientID, row[0], '', row[1], row[2], row[3], row[4], 1 # 补充其他列参数 ) batch.append(row_args) if len(batch) >= batch_size: process_batch(batch) batch = [] if batch: process_batch(batch) def ReadMetaDataFile(srcFile, fileDelimiter, clientID): functionResponse = '' try: # 读取文件并分割为多个块(示例按行数分割) with open(srcFile, 'r', buffering=1024*1024) as delimitedFile: fileReader = csv.reader(delimitedFile, delimiter=fileDelimiter) next(fileReader) # 跳过表头 rows = list(fileReader) chunk_size = len(rows) // 4 + 1 chunks = [rows[i:i+chunk_size] for i in range(0, len(rows), chunk_size)] # 用线程池处理各块 with ThreadPoolExecutor(max_workers=4) as executor: for chunk in chunks: executor.submit(process_file_chunk, chunk, clientID) functionResponse = "MetaData File Loaded Successfully" except Exception as err: functionResponse = f"Exception Found: {str(err)}" return functionResponse
内容的提问来源于stack exchange,提问作者arbennett4
相关产品推荐
相关产品推荐

