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

循环执行SQLAlchemy插入语句时固定触发MySQL连接丢失问题求助

解决循环执行INSERT...SELECT时固定迭代次数触发MySQL连接丢失问题

问题场景

循环执行INSERT INTO table_1 SELECT * FROM table_2 WHERE ... ON DUPLICATE KEY UPDATE语句时,固定在第150次迭代触发连接丢失异常:

sqlalchemy.exc.OperationalError: (pymysql.err.OperationalError) (2013, 'Lost connection to MySQL server during query')

已排除以下可能性:

  • 循环总耗时仅约5秒,不存在超时问题
  • 同逻辑Java实现无异常
  • 总数据量仅5000行(不足300KB),数据负载极低
  • 每次迭代后已执行db.session.commit()
  • 独立脚本使用单独Engine和SessionMaker仍复现问题

测试代码片段:

normal_date_pattern = '%Y-%m-%d'
day_start_hour_format_pattern = '%Y-%m-%d %H:00:00'
day_end_hour_format_pattern = '%Y-%m-%d %H:59:59'

while actual_end_time > start_time:
    start_current_hour_str = start_time.strftime(self.day_start_hour_format_pattern)
    end_current_hour_str = start_time.strftime(self.day_end_hour_format_pattern)
    # log when meet 0'oclock
    if start_current_hour_str.endswith('00:00:00'):
        current_app.logger.info(f'processing date {start_time.strftime(self.normal_date_pattern)}')
    # insert into db
    db.session.execute(
        text(
            """
                INSERT INTO table_1
                SELECT * FROM table_2
                WHERE table_2.start_time>=:start_time
                AND table_2.end_time<=:end_time
                ON DUPLICATE KEY UPDATE
                    snapshot_date = snapshot_date,
                    snapshot_date_loc = VALUES(snapshot_date_loc)
            """
        ),
        params={
            'start_time': start_current_hour_str,
            'end_time': end_current_hour_str
        }
    )
    db.session.commit()

    # add one hour to process next hour
    start_time = start_time + timedelta(hours=1)

解决方案

1. 检查并调整MySQL的max_allowed_packet参数

虽然单条数据量小,但高频执行的语句可能累积数据包大小,触及MySQL默认的max_allowed_packet阈值(通常为4MB)。临时调整测试:

SET GLOBAL max_allowed_packet=67108864; -- 设置为64MB,临时生效

若问题解决,需将该参数永久写入MySQL配置文件(如my.cnf):

[mysqld]
max_allowed_packet=64M

2. 减少执行次数,改为批量时间段处理

没必要按小时单条执行,合并多个小时的时间区间,大幅降低语句执行次数:

batch_hours = 10  # 每次处理10小时数据
while actual_end_time > start_time:
    batch_end = min(start_time + timedelta(hours=batch_hours), actual_end_time)
    start_current_str = start_time.strftime(day_start_hour_format_pattern)
    end_current_str = batch_end.strftime(day_end_hour_format_pattern)
    
    db.session.execute(
        text("""
            INSERT INTO table_1
            SELECT * FROM table_2
            WHERE table_2.start_time >= :start_time
              AND table_2.end_time <= :end_time
            ON DUPLICATE KEY UPDATE
                snapshot_date = snapshot_date,
                snapshot_date_loc = VALUES(snapshot_date_loc)
        """),
        params={'start_time': start_current_str, 'end_time': end_current_str}
    )
    db.session.commit()
    start_time = batch_end

3. 每次迭代后显式重置Session

即使执行了commit,单个Session长时间持有连接可能累积异常,尝试每次迭代后关闭并重建Session:

while actual_end_time > start_time:
    # ... 生成时间参数 ...
    db.session.execute(/* SQL语句 */)
    db.session.commit()
    db.session.close()  # 释放连接回池
    db.session = Session()  # 重新创建Session(需确保Session是sessionmaker实例)
    
    start_time = start_time + timedelta(hours=1)

或改用上下文管理器自动管理Session生命周期:

while actual_end_time > start_time:
    # ... 生成时间参数 ...
    with Session() as session:
        session.execute(/* SQL语句 */)
        session.commit()
    
    start_time = start_time + timedelta(hours=1)

4. 升级SQLAlchemy和pymysql版本

版本兼容性问题可能导致高频操作时连接异常,升级到最新稳定版:

pip install --upgrade sqlalchemy pymysql

5. 查看MySQL日志定位根因

开启MySQL慢查询日志和错误日志,查看第150次执行时的服务器端细节:

  • 开启慢查询日志:在my.cnf中添加slow_query_log = 1、long_query_time = 0(记录所有查询)
  • 查看错误日志位置:SHOW VARIABLES LIKE 'log_error';
    通过日志确认MySQL是否主动断开连接,以及具体触发原因。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 15:27:47