如何在BQ表值变更时触发GCP VM中Scala/Python代码执行ETL
实现方案
1. 设计BigQuery任务表结构
先建一张BQ表管理ETL任务,核心字段覆盖任务识别、状态追踪和执行信息:
CREATE TABLE `your-project.your-dataset.etl_tasks` ( task_id STRING NOT NULL, -- 唯一任务ID,用UUID生成 task_type STRING NOT NULL, -- 任务类型:'scala' 或 'python' script_path STRING NOT NULL, -- ETL脚本路径:VM本地路径或GCS路径都可 status STRING NOT NULL DEFAULT 'pending', -- 状态:pending/processing/success/failed create_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP(), update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP(), log_path STRING, -- 日志存储路径(比如GCS) retry_count INT64 DEFAULT 0 -- 重试次数 ) PRIMARY KEY(task_id);
主键避免重复任务,状态字段用来筛选待执行任务,防止重复处理。
2. 编写VM上的轮询调度脚本
用Python写核心调度逻辑(处理多进程和BQ交互更灵活),核心逻辑是:定时查询BQ的待执行任务,原子标记为执行中,并行启动ETL任务,最后更新任务状态。
核心代码示例
import time import subprocess from google.cloud import bigquery from concurrent.futures import ProcessPoolExecutor # 初始化BQ客户端 client = bigquery.Client() PROJECT_ID = "your-project" DATASET_ID = "your-dataset" TABLE_ID = "etl_tasks" TABLE_FULL_ID = f"{PROJECT_ID}.{DATASET_ID}.{TABLE_ID}" # 轮询间隔(单位:秒,比如60秒一次) POLL_INTERVAL = 60 def execute_etl_task(task): """执行单个ETL任务""" task_id = task["task_id"] task_type = task["task_type"] script_path = task["script_path"] log_path = f"gs://your-log-bucket/etl_logs/{task_id}.log" try: # 根据任务类型调用对应执行命令 if task_type == "python": cmd = ["python3", script_path] elif task_type == "scala": # 假设Scala脚本已打包成jar,用spark-submit运行 cmd = ["spark-submit", "--class", "com.your.etl.MainClass", script_path] else: raise ValueError(f"Unknown task type: {task_type}") # 执行命令并暂存日志到本地 with open(f"/tmp/{task_id}.log", "w") as f: subprocess.run(cmd, check=True, stdout=f, stderr=f) # 上传日志到GCS subprocess.run(["gsutil", "cp", f"/tmp/{task_id}.log", log_path], check=True) # 更新任务状态为success update_query = f""" UPDATE `{TABLE_FULL_ID}` SET status = 'success', update_time = CURRENT_TIMESTAMP(), log_path = '{log_path}' WHERE task_id = '{task_id}' """ client.query(update_query).result() except Exception as e: # 更新任务状态为failed,重试次数+1 update_query = f""" UPDATE `{TABLE_FULL_ID}` SET status = 'failed', update_time = CURRENT_TIMESTAMP(), retry_count = retry_count + 1 WHERE task_id = '{task_id}' """ client.query(update_query).result() print(f"Task {task_id} failed: {str(e)}") def poll_bq_tasks(): """轮询BQ表,处理待执行任务""" while True: # 原子操作:先把pending任务标记为processing,避免多进程重复处理 lock_query = f""" UPDATE `{TABLE_FULL_ID}` SET status = 'processing', update_time = CURRENT_TIMESTAMP() WHERE status = 'pending' RETURNING task_id, task_type, script_path """ query_job = client.query(lock_query) tasks = [dict(row) for row in query_job.result()] if tasks: print(f"发现{len(tasks)}个待执行任务,开始执行...") # 用进程池并行执行,max_workers根据VM配置调整 with ProcessPoolExecutor(max_workers=4) as executor: executor.map(execute_etl_task, tasks) else: print("无待执行任务,等待下一轮轮询...") time.sleep(POLL_INTERVAL) if __name__ == "__main__": poll_bq_tasks()
3. 配置VM环境
- 安装依赖:执行
pip install google-cloud-bigquery,确保gsutil、python3、spark-submit(如果用Scala)已安装。 - 权限配置:给VM的服务账号添加BigQuery读写权限、GCS读写权限,可通过IAM添加
BigQuery Data Editor和Storage Object Admin角色。
4. 持久化轮询服务
把调度脚本做成systemd服务,确保VM重启后自动运行:
- 创建服务文件
/etc/systemd/system/etl-scheduler.service:
[Unit] Description=ETL Task Scheduler After=network.target [Service] User=your-vm-user ExecStart=/usr/bin/python3 /path/to/your/scheduler_script.py Restart=always RestartSec=10 [Install] WantedBy=multi-user.target
- 启动并启用服务:
sudo systemctl daemon-reload sudo systemctl start etl-scheduler sudo systemctl enable etl-scheduler
5. 额外优化点
- 轮询间隔:根据任务频率调整,避免过于频繁查询BQ增加成本。
- 重试机制:可在调度脚本中判断
retry_count,若小于3则将failed状态改回pending重新执行。 - 资源限制:根据VM的CPU和内存配置,调整进程池
max_workers,避免资源耗尽。 - 任务去重:利用BQ表主键
task_id,避免插入重复任务。
内容的提问来源于stack exchange,提问作者Khilesh Chauhan
相关产品推荐
相关产品推荐

