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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 14:20:17