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

