BigQuery Python API会话事务ROLLBACK后数据未回滚问题排查
BigQuery Python API 跨查询事务回滚失效问题
问题表现
基于BigQuery会话实现多语句事务时,异常触发ROLLBACK TRANSACTION执行成功后,目标表内依然存在已写入的数据,回滚逻辑未生效。
复现代码
会话上下文管理器实现:
from google.cloud import bigquery class BigquerySession: """封装BigQuery会话的上下文管理器""" def __init__(self, bqclient: bigquery.Client, bqlocation: str = "EU") -> None: self._bigquery_client = bqclient self._location = bqlocation self._session_id = None def __enter__(self) -> str: """初始化会话并返回会话ID""" job = self._bigquery_client.query( "SELECT 1;", job_config=bigquery.QueryJobConfig(create_session=True), location=self._location, ) self._session_id = job.session_info.session_id job.result() return self._session_id def __exit__(self, exc_type, exc_value, traceback): """退出时清理会话""" if exc_type: print("事务执行失败,触发回滚") job = self._bigquery_client.query( "ROLLBACK TRANSACTION;", job_config=bigquery.QueryJobConfig( create_session=False, connection_properties=[ bigquery.query.ConnectionProperty(key="session_id", value=self._session_id) ], ), location=self._location, ) job.result() if self._session_id: # 强制终止会话,避免表锁残留 job = self._bigquery_client.query( "CALL BQ.ABORT_SESSION();", job_config=bigquery.QueryJobConfig( create_session=False, connection_properties=[ bigquery.query.ConnectionProperty( key="session_id", value=self._session_id ) ], ), location=self._location, ) job.result() return False
事务测试逻辑:
# 开启事务 job = self.client.query( "BEGIN TRANSACTION;", job_config=bigquery.QueryJobConfig( create_session=False, connection_properties=[ bigquery.query.ConnectionProperty(key="session_id", value=session_id) ] ), location=self.dataset.location, ) job.result() # 执行聚合查询并写入目标表 job = self.client.query( aggregation_query, job_config=bigquery.QueryJobConfig( create_session=False, connection_properties=[ bigquery.query.ConnectionProperty(key="session_id", value=session_id) ], destination=f"{self.dataset.project}.{self.dataset.dataset_id}.{table_name}", create_disposition="CREATE_NEVER", write_disposition="WRITE_APPEND" ), location=self.dataset.location, ) print(job.result()) # 主动抛出异常,跳过后续提交逻辑 raise KeyboardInterrupt # 提交事务(不会被执行) job = self.client.query( "COMMIT TRANSACTION;", job_config=bigquery.QueryJobConfig( create_session=False, connection_properties=[ bigquery.query.ConnectionProperty(key="session_id", value=session_id) ], ), location=self.dataset.location, ) job.result()
根因
第一个猜想成立,问题和API预览阶段缺陷无关:
- 配置
QueryJobConfig的destination、write_disposition参数触发的结果落表,是BigQuery作业调度层的独立逻辑,在SQL语句执行完成后触发,不属于事务块内的DML操作,完全不受事务回滚规则约束。 - BigQuery事务仅能管控SQL文本内部的DML操作,即直接写在SQL中的
INSERT/UPDATE/DELETE/MERGE语句,这类操作产生的变更才会在未提交时被ROLLBACK撤销。当前测试逻辑中执行的SQL是纯SELECT聚合语句,落表动作是作业配置带来的额外操作,不在事务边界内,因此回滚无法撤销这部分写入。
修复方案
- 事务内的表写入不要依赖作业配置的落表参数实现,直接将写入逻辑写入SQL文本,使用标准DML语法,示例:
这类DML操作完全运行在事务上下文内,回滚时会完整撤销所有未提交的写入。INSERT INTO `项目ID.数据集ID.目标表名` -- 替换为原有聚合查询逻辑 SELECT dimension, SUM(metric) FROM source_table GROUP BY dimension - 移除上下文管理器中手动执行
ROLLBACK TRANSACTION的逻辑:BigQuery会话自带未提交事务自动回滚机制,只要会话终止前没有执行COMMIT,所有事务内变更都会自动回滚,手动调用ROLLBACK在部分会话异常场景下反而会触发报错,保留BQ.ABORT_SESSION()做会话清理即可。 - 如果业务场景必须使用作业结果直接落表的能力,不要将这类操作放入事务块。单作业的结果落表本身是原子性的,要么完整写入成功,要么作业失败无任何写入,不需要额外事务做一致性保障。
内容的提问来源于stack exchange,提问作者Tonca
相关产品推荐
相关产品推荐

