多线程(gevent池)下pymssql库出现SQL事务错误求助
问题:gevent池规模增大时出现事务提交错误
错误信息:The COMMIT TRANSACTION request has no corresponding BEGIN TRANSACTION.
gevent池规模越大(如20及以上)错误越频繁,10左右时几乎无错误。原本怀疑是其他包含cursor.execute("BEGIN TRAN....")的代码导致,但实际这类场景极少。
代码场景
- 主函数通过gevent池启动任务:
pool.spawn(collector, target_table, source_table, columns_to_copy, sql_filter)
- 任务函数
collector实例化自定义MySQLClass创建数据库连接:
def collector(target_table, source_table, columns_to_copy, sql_filter): mydb = MySQLClass(host=SQL_HOST, database_name='mydb', user=myuser, password=mypw) # ... 其他业务逻辑 mydb.sql_delete(table, sql_where_filter)
class MySQLClass(object): def __init__(self, host, database_name, user, password): self.db = pymssql.connect( host=host, database=database_name, user=user, password=password ) self.cursor = self.db.cursor()
sql_delete方法执行删除后提交事务:
def sql_delete(self, table, sql_filter=""): self.cursor.execute("DELETE FROM " + table + " " + sql_filter ) self.db.commit() return
解决思路
1. 显式控制事务生命周期
pymssql默认关闭自动提交,但gevent协程切换可能干扰连接的事务状态。建议显式开启事务后再执行操作,避免状态混乱:
def sql_delete(self, table, sql_filter=""): try: self.db.begin() # 显式启动事务 self.cursor.execute("DELETE FROM " + table + " " + sql_filter ) self.db.commit() except Exception as e: self.db.rollback() # 异常时回滚 raise e
也可在初始化连接时开启自动提交,跳过手动事务管理:
self.db = pymssql.connect( host=host, database=database_name, user=user, password=password, autocommit=True # 开启自动提交 )
2. 排查协程连接共享问题
gevent的猴子补丁可能影响pymssql底层socket处理,导致协程意外共享连接。可采取以下措施:
- 确保每个协程实例化的
MySQLClass对应独立的数据库连接,避免连接池复用冲突 - 替换为协程友好的数据库库(如
asyncpymssql)替代pymssql,适配gevent的异步模型
3. 增加连接有效性检查
当连接因网络波动或数据库端断开时,协程可能复用失效连接,导致事务状态异常。建议执行操作前检查连接状态:
def _check_connection(self): try: self.cursor.execute("SELECT 1") except Exception: # 连接失效,重建连接 self.db.close() self.db = pymssql.connect( host=self.host, database=self.database_name, user=self.user, password=self.password ) self.cursor = self.db.cursor() def sql_delete(self, table, sql_filter=""): self._check_connection() self.cursor.execute("DELETE FROM " + table + " " + sql_filter ) self.db.commit()
4. 修复SQL拼接的潜在风险
当前SQL拼接方式存在注入风险,若sql_filter包含意外的事务语句(如误加COMMIT),会直接打乱事务状态。改用参数化查询规避:
def sql_delete(self, table, sql_filter="", params=None): sql = f"DELETE FROM {table}" if sql_filter: sql += " " + sql_filter self.cursor.execute(sql, params or ()) self.db.commit()
内容的提问来源于stack exchange,提问作者Fabri Ba
相关产品推荐
相关产品推荐

