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

MySQL数据库轮询无法获取新结果的问题求助

解决MySQL长连接轮询无法获取新数据的问题

问题场景

需同步遗留系统生成的MySQL表数据至另一服务器,新增posted_in_prodmon字段(默认值0),轮询逻辑为:查询该字段为0的记录,同步完成后将其更新为1,无新记录则休眠10秒重试。新数据由其他进程写入(已通过Adminer确认提交)。

环境:Python3.9 + mysql-connector-python==8.0.33,MySQL 8.0.23

遇到的问题:首次查询能获取所有预期记录,但后续循环无法读取新数据,仅重启程序或每次关闭连接重连有效,重连效率过低,需更优解决方案。

原因分析

MySQL InnoDB引擎默认事务隔离级别为REPEATABLE READ(可重复读),在同一事务内,所有查询都会基于事务启动时的数据快照,无法读取其他事务提交的新数据。

你的代码中,mysql-connector默认关闭自动提交(autocommit=False),首次查询后事务未显式结束,后续所有查询都处于同一事务中,持续使用初始快照,因此看不到其他进程写入的新数据。而重连会创建新连接并启动新事务,自然能读取最新数据。

解决方案

方案1:开启自动提交

建立数据库连接时设置autocommit=True,让每个查询/操作自动提交事务,每次轮询都基于最新数据快照查询。

修改Mysql_DB类的is_connected方法:

def is_connected(self):
    if self.connection:
        if self.connection.is_connected():
            return True
    try:
        self.logger.info(f'Not connected to mysql server... reconnecting')  
        # 添加autocommit=True参数
        self.connection= mysql.connector.connect(**self.dbconfig, autocommit=True)
        return True

    except (mysql.connector.Error, IOError) as err:
        self.logger.error(f'Mysql connection failed: {err}')
        return False                

方案2:手动结束事务

若不需要全局自动提交,可在每次轮询结束后手动提交或回滚事务(即使无修改操作,SELECT也会启动事务),确保下次查询在新事务中执行:

修改循环逻辑,替换source.connection.close()为事务结束操作:

while True:
    records = source.get_records()
    for result in records:
        id = result[0]
        data = result[1]
        data_list = data.split("\t",3)
        
        created_at = f'{data_list[0]} {data_list[1]}'
        created_at = created_at.replace('_', '-')
        part = data_list[2]
        
        ts = int(datetime.timestamp(datetime.fromisoformat(created_at)))
        target.insert(ts, part)
        rowcount = source.update_record(id, created_at)
        if not rowcount:
            raise Exception(f'Failed to update id {id}')

    # 结束当前事务,释放数据快照
    source.connection.commit()  # 或使用source.connection.rollback(),效果一致
    sleep(10)

同时优化get_records方法,直接返回结果列表(避免游标资源泄漏):

class Source(Mysql_DB):        
    def get_records(self):
        if not self.is_connected():
            return []
        sql  = 'SELECT id, inspection_data '
        sql += 'FROM `1730_Vantage` '
        sql += 'WHERE part_fail = 1 '
        sql += 'AND posted_in_GFxPRoduction = 0 '
        sql += 'AND created_at > "2023-06-01" '
        sql += 'ORDER BY created_at DESC; '

        with self.connection.cursor(buffered=True) as cursor:
            cursor.execute(sql)
            return cursor.fetchall()  # 返回所有结果的列表

方案3:调整事务隔离级别

若业务逻辑允许,可将连接的事务隔离级别改为READ COMMITTED(读已提交),这样同一事务内的后续查询能读取其他事务提交的新数据。

在连接后设置隔离级别:

self.connection= mysql.connector.connect(**self.dbconfig)
self.connection.start_transaction(isolation_level='READ COMMITTED')

注意:此方案需确保业务兼容读已提交的特性(可能出现不可重复读)。

关键注意事项

  • 避免长时间持有未结束的事务,不仅会导致数据快照无法更新,还会占用数据库连接资源
  • 游标使用后需及时关闭,推荐用with语句自动管理游标生命周期
  • 优先选择自动提交或手动结束事务的方案,效率远高于每次重连

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 07:55:54