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

