使用Python或BQ CLI自动化BigQuery表存在性检查及运维操作
自动化BigQuery表同步清理脚本(Python + BQ CLI)
前提条件
- 已安装并配置BQ CLI(完成
gcloud auth application-default login认证) - 准备好包含表名的文本文件(每行一个表名,例如
tables.txt) - 拥有目标项目的BigQuery操作权限(Data Editor及以上)
脚本实现
import subprocess import sys # 配置项目和数据集信息,根据实际情况修改 PROJECTS = { "dev": "my-project-dev", "test": "my-project-test", "prod": "my-project-prod" } DATASETS = ["my_dataset_hist", "my_dataset_curr", "my_dataset_stg"] STAGING_SUFFIX = "_stg" def run_bq_command(cmd): """执行BQ CLI命令并返回结果""" try: result = subprocess.run( cmd, shell=True, check=True, capture_output=True, text=True ) print(f"成功执行命令: {cmd}") return True except subprocess.CalledProcessError as e: print(f"命令执行失败: {cmd}\n错误信息: {e.stderr}", file=sys.stderr) return False def check_table_exists(project, dataset, table): """检查表是否存在于指定项目和数据集""" cmd = f"bq show --project_id {project} {dataset}.{table} > /dev/null 2>&1" return subprocess.run(cmd, shell=True).returncode == 0 def process_table(table_name): """处理单个表的清理/同步逻辑""" print(f"\n=== 开始处理表: {table_name} ===") # 收集各环境下的表存在情况 prod_exists = any(check_table_exists(PROJECTS["prod"], ds, table_name) for ds in DATASETS) dev_table_locations = [(ds, table_name) for ds in DATASETS if check_table_exists(PROJECTS["dev"], ds, table_name)] test_table_locations = [(ds, table_name) for ds in DATASETS if check_table_exists(PROJECTS["test"], ds, table_name)] # 检查是否属于staging数据集(规则4优先级最高) staging_datasets = [ds for ds in DATASETS if ds.endswith(STAGING_SUFFIX)] is_staging_table = any( check_table_exists(proj, ds, table_name) for proj in PROJECTS.values() for ds in staging_datasets ) if is_staging_table: print("检测到该表属于staging数据集,执行截断操作") # 截断Dev环境的表 for ds, tbl in dev_table_locations: truncate_cmd = f"bq query --project_id {PROJECTS['dev']} --use_legacy_sql=false 'TRUNCATE TABLE `{PROJECTS['dev']}.{ds}.{tbl}`'" run_bq_command(truncate_cmd) # 截断Test环境的表 for ds, tbl in test_table_locations: truncate_cmd = f"bq query --project_id {PROJECTS['test']} --use_legacy_sql=false 'TRUNCATE TABLE `{PROJECTS['test']}.{ds}.{tbl}`'" run_bq_command(truncate_cmd) return # 处理非staging表的逻辑 if prod_exists: print("Prod环境存在该表,执行Dev/Test表截断操作") # 截断Dev环境的表 for ds, tbl in dev_table_locations: truncate_cmd = f"bq query --project_id {PROJECTS['dev']} --use_legacy_sql=false 'TRUNCATE TABLE `{PROJECTS['dev']}.{ds}.{tbl}`'" run_bq_command(truncate_cmd) # 截断Test环境的表 for ds, tbl in test_table_locations: truncate_cmd = f"bq query --project_id {PROJECTS['test']} --use_legacy_sql=false 'TRUNCATE TABLE `{PROJECTS['test']}.{ds}.{tbl}`'" run_bq_command(truncate_cmd) else: print("Prod环境不存在该表,执行Dev/Test表删除操作") # 删除Dev环境的表 for ds, tbl in dev_table_locations: drop_cmd = f"bq rm --project_id {PROJECTS['dev']} -f -t `{PROJECTS['dev']}.{ds}.{tbl}`" run_bq_command(drop_cmd) # 删除Test环境的表 for ds, tbl in test_table_locations: drop_cmd = f"bq rm --project_id {PROJECTS['test']} -f -t `{PROJECTS['test']}.{ds}.{tbl}`" run_bq_command(drop_cmd) if __name__ == "__main__": if len(sys.argv) != 2: print("用法: python bq_table_cleanup.py <表名文件路径>") sys.exit(1) table_file = sys.argv[1] try: with open(table_file, "r") as f: tables = [line.strip() for line in f if line.strip()] except FileNotFoundError: print(f"错误:未找到文件 {table_file}", file=sys.stderr) sys.exit(1) for table in tables: process_table(table) print("\n=== 所有表处理完成 ===")
使用说明
- 修改脚本顶部的
PROJECTS和DATASETS字典/列表,匹配你的实际环境 - 准备好包含表名的文本文件(每行一个表名,无多余空格)
- 执行脚本:
python bq_table_cleanup.py tables.txt
关键逻辑说明
- 规则4优先级最高:只要表存在于任何以
*_stg结尾的数据集,无论Prod是否存在,仅执行Dev/Test表的截断操作 - 规则2:Prod存在该表时,仅清空Dev/Test的对应表内容,不删除
- 规则3:Prod不存在该表时,直接删除Dev/Test的对应表
- 脚本会自动遍历所有指定数据集,查找目标表的存在情况
内容的提问来源于stack exchange,提问作者marie20
相关产品推荐
相关产品推荐

