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

Flask+APScheduler执行远程SQL文件时遇Commands out of sync错误

问题原因分析

1. 未处理multi=True生成的所有结果集

使用cursor.execute(sql, multi=True)执行多语句SQL时,每个SQL命令都会生成一个结果集。如果不遍历并处理完所有结果集,MySQL连接会处于"命令未同步"的状态,后续的commit操作就会触发Commands out of sync错误。

2. 全局数据库连接的线程安全冲突

APScheduler默认采用多线程执行定时任务,而你复用了全局的db连接对象。MySQL连接本身不是线程安全的,多线程同时操作同一个连接会导致连接状态混乱,进而引发同步错误。

3. 错误的连接关闭时机

在finally块中调用db.close()关闭了全局连接,后续任务执行时会复用已经失效的连接,这也会导致连接状态异常,触发同步错误。


解决方法

方法1:处理所有multi=True生成的结果集

在执行多语句SQL后,必须遍历所有结果集,让MySQL连接回到同步状态:

def execute_file(sql_file):
    if sql_file:
        try:
            with open(sql_file, 'r') as sql_f:
                sql = sql_f.read()

            cursor = db.cursor(dictionary=True)
            # 执行多语句并获取结果集迭代器
            results = cursor.execute(sql, multi=True)
            # 遍历所有结果集,完成命令同步
            for result in results:
                if result.with_rows:
                    # 对于有返回结果的语句(如SELECT),读取结果(这里不需要可以直接跳过)
                    result.fetchall()
            db.commit()
        # 后续异常处理、关闭逻辑不变...

方法2:每个任务创建独立的数据库连接

放弃全局连接复用,改为每个任务执行时创建新的连接,彻底避免多线程冲突:

from mysql.connector import errors
import mysql.connector

# 替换成你的数据库连接配置
def get_db_connection():
    return mysql.connector.connect(
        host="your_remote_host",
        user="db_user",
        password="db_password",
        # 注意:初始连接可以不指定数据库,因为SQL里有USE语句
    )

def execute_file(sql_file):
    if sql_file:
        db = None
        cursor = None
        try:
            # 每次任务创建独立连接
            db = get_db_connection()
            with open(sql_file, 'r') as sql_f:
                sql = sql_f.read()

            cursor = db.cursor(dictionary=True)
            results = cursor.execute(sql, multi=True)
            # 遍历所有结果集
            for result in results:
                if result.with_rows:
                    result.fetchall()
            db.commit()
        except errors.Error as e:
            print(f'[ ERROR ]: Error while Executing {str(e)}')
            print(f'[ ERROR ]: Rolling back ...')
            if db:
                db.rollback()
            return False, str(e)
        except Exception as e:
            print(f'[ ERROR ]: Error while Executing {str(e)}')
            print(f'[ ERROR ]: Rolling back ...')
            if db:
                db.rollback()
            return False, str(e)
        finally:     
            if cursor:
                cursor.close() 
            if db:
                db.close()
        print(f'[ INFO ]: Transaction committed....')
        return True, "Transaction committed."
    else:
        print(f'[ ERROR ]: Sql File not found ...')
        return False, "Sql File not found"

方法3:临时改用单线程执行器(快速缓解)

如果暂时不想修改连接逻辑,可以将APScheduler配置为单线程执行,避免多线程共享连接的问题:

from execute_sql import execute_file
from flask_apscheduler import APScheduler

scheduler = APScheduler()
# 配置单线程执行器
scheduler.init_app(app)
scheduler.executors = {
    'default': {'type': 'threadpool', 'max_workers': 1}
}
scheduler.start()

files = ["path/to/file1.sql", "path/to/file2.sql"]
for idx, file_path in enumerate(files):
    job_id = f"deploy_job_{idx}"
    scheduler.add_job(func=execute_file, id=job_id, trigger='cron', coalesce=True, args=[file_path])

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 02:36:29