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

如何在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重启后自动运行:

  1. 创建服务文件/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
  1. 启动并启用服务:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 13:10:24