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

SQLAlchemy中PendingRollbackError问题:MySQL写入中断后无法回滚

解决SQLAlchemy + pandas to_sql的PendingRollbackError问题

问题描述

使用SQLAlchemy和pandas的to_sql()方法将CSV文件写入MySQL表时程序意外退出,导致存在无法回滚的无效事务,触发PendingRollbackError错误,提示信息:

Can't reconnect until invalid transaction is rolled back. Please rollback() fully before proceeding

尝试通过try-except块执行回滚操作,但未解决问题。

报错栈信息

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\pandas\core\generic.py:2878, in NDFrame.to_sql(self, name, con, schema, if_exists, index, index_label, chunksize, dtype, method)
   2713 """
   2714 Write records stored in a DataFrame to a SQL database.
   2715 
   (...)
   2874 [(1,), (None,), (2,)]
   2875 """  # noqa:E501
   2876 from pandas.io import sql
-> 2878 return sql.to_sql(
   2879     self,
   2880     name,
   2881     con,
   2882     schema=schema,
   2883     if_exists=if_exists,
   2884     index=index,
   2885     index_label=index_label,
   2886     chunksize=chunksize,
   2887     dtype=dtype,
   2888     method=method,
   2889 )

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\pandas\io\sql.py:768, in to_sql(frame, name, con, schema, if_exists, index, index_label, chunksize, dtype, method, engine, **engine_kwargs)
    763 elif not isinstance(frame, DataFrame):
    764     raise NotImplementedError(
    765         "'frame' argument should be either a Series or a DataFrame"
    766     )
-> 768 with pandasSQL_builder(con, schema=schema, need_transaction=True) as pandas_sql:
    769     return pandas_sql.to_sql(
    770         frame,
    771         name,
   (...)
    780         **engine_kwargs,
    781     )

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\pandas\io\sql.py:1548, in SQLDatabase.__exit__(self, *args)
   1546 def __exit__(self, *args) -> None:
   1547     if not self.returns_generator:
-> 1548         self.exit_stack.close()

File ~\AppData\Local\Programs\Python\Python311\Lib\contextlib.py:597, in ExitStack.close(self)
    595 def close(self):
    596     """Immediately unwind the context stack."""
-> 597     self.__exit__(None, None, None)

File ~\AppData\Local\Programs\Python\Python311\Lib\contextlib.py:589, in ExitStack.__exit__(self, *exc_details)
    585 try:
    586     # bare "raise exc_details[1]" replaces our carefully
    587     # set-up context
    588     fixed_ctx = exc_details[1].__context__
-> 589     raise exc_details[1]
    590 except BaseException:
    591     exc_details[1].__context__ = fixed_ctx

File ~\AppData\Local\Programs\Python\Python311\Lib\contextlib.py:574, in ExitStack.__exit__(self, *exc_details)
    572 assert is_sync
    573 try:
-> 574     if cb(*exc_details):
    575         suppressed_exc = True
    576         pending_raise = False

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\util.py:146, in TransactionalContext.__exit__(self, type_, value, traceback)
    144     self.commit()
    145 except:
-> 146     with util.safe_reraise():
    147         if self._rollback_can_be_called():
    148             self.rollback()

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\util\langhelpers.py:147, in safe_reraise.__exit__(self, type_, value, traceback)
    145     assert exc_value is not None
    146     self._exc_info = None  # remove potential circular references
-> 147     raise exc_value.with_traceback(exc_tb)
    148 else:
    149     self._exc_info = None  # remove potential circular references

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\util.py:144, in TransactionalContext.__exit__(self, type_, value, traceback)
    142 if type_ is None and self._transaction_is_active():
    143     try:
-> 144         self.commit()
    145     except:
    146         with util.safe_reraise():

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:2615, in Transaction.commit(self)
   2599 """Commit this :class:`.Transaction`.
   2600 
   2601 The implementation of this may vary based on the type of transaction in
   (...)
   2612 
   2613 """
   2614 try:
-> 2615     self._do_commit()
   2616 finally:
   2617     assert not self.is_active

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:2720, in RootTransaction._do_commit(self)
   2717 assert self.connection._transaction is self
   2719 try:
-> 2720     self._connection_commit_impl()
   2721 finally:
   2722     # whether or not commit succeeds, cancel any
   2723     # nested transactions, make this transaction "inactive"
   2724     # and remove it as a reset agent
   2725     if self.connection._nested_transaction:

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:2691, in RootTransaction._connection_commit_impl(self)
   2690 def _connection_commit_impl(self) -> None:
-> 2691     self.connection._commit_impl()

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:1134, in Connection._commit_impl(self)
   1132     self.engine.dialect.do_commit(self.connection)
   1133 except BaseException as e:
-> 1134     self._handle_dbapi_exception(e, None, None, None, None)

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:2342, in Connection._handle_dbapi_exception(self, e, statement, parameters, cursor, context, is_sub_exec)
   2340     else:
   2341         assert exc_info[1] is not None
-> 2342         raise exc_info[1].with_traceback(exc_info[2])
   2343 finally:
   2344     del self._reentrant_error

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:1132, in Connection._commit_impl(self)
   1130         self._log_info("COMMIT")
   1131 try:
-> 1132     self.engine.dialect.do_commit(self.connection)
   1133 except BaseException as e:
   1134     self._handle_dbapi_exception(e, None, None, None, None)

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:573, in Connection.connection(self)
    571 if self._dbapi_connection is None:
    572     try:
-> 573         return self._revalidate_connection()
    574     except (exc.PendingRollbackError, exc.ResourceClosedError):
    575         raise

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:665, in Connection._revalidate_connection(self)
    663 if self.__can_reconnect and self.invalidated:
    664     if self._transaction is not None:
-> 665         self._invalid_transaction()
    666     self._dbapi_connection = self.engine.raw_connection()
    667     return self._dbapi_connection

File ~\AppData\Local\Programs\Python\Python311\Lib\site-packages\sqlalchemy\engine\base.py:655, in Connection._invalid_transaction(self)
    654 def _invalid_transaction(self) -> NoReturn:
-> 655     raise exc.PendingRollbackError(
    656         "Can't reconnect until invalid %stransaction is rolled "
    657         "back.  Please rollback() fully before proceeding"
    658         % ("savepoint " if self._nested_transaction is not None else ""),
    659         code="8s2b",
    660     )

解决方案

1. 数据库层面手动清理无效事务

登录MySQL控制台,执行以下步骤:

  • 查看当前活跃事务:
    SELECT * FROM INFORMATION_SCHEMA.INNODB_TRX;
    
  • 找到对应的未提交事务,记录其trx_mysql_thread_id,然后终止该会话:
    KILL [trx_mysql_thread_id];
    
  • 若为当前会话产生的问题,直接执行回滚:
    ROLLBACK;
    

2. 重置SQLAlchemy连接并显式回滚

不要复用之前的连接实例,创建新的引擎和连接,强制回滚无效事务:

from sqlalchemy import create_engine
from sqlalchemy.exc import PendingRollbackError

# 替换为你的数据库连接信息
engine = create_engine('mysql+pymysql://username:password@host:port/dbname')

with engine.connect() as conn:
    try:
        conn.rollback()
        print("事务回滚成功")
    except PendingRollbackError as e:
        print(f"回滚失败: {e}")
        engine.dispose()

3. 手动控制pandas to_sql的事务

使用SQLAlchemy的上下文管理器自动处理事务,确保异常时正确回滚:

from sqlalchemy import create_engine
from sqlalchemy.exc import SQLAlchemyError
import pandas as pd

engine = create_engine('mysql+pymysql://username:password@host:port/dbname')
df = pd.read_csv('your_file.csv')

# engine.begin()会自动处理:成功提交,异常回滚
with engine.begin() as conn:
    try:
        df.to_sql(
            name='target_table',
            con=conn,
            if_exists='append',
            chunksize=1000,
            index=False
        )
        print("数据写入成功")
    except SQLAlchemyError as e:
        print(f"写入失败: {e}")

4. 清理SQLAlchemy连接池

如果使用了连接池,池内的连接可能处于无效事务状态,直接销毁连接池:

engine.dispose()

之后重新创建连接即可正常使用。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 10:54:54