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

如何在GCP Airflow中便捷监控BQ Load任务的Job ID状态?

实现BigQuery Load任务状态监控的简便方案

一、用BigQuery Python SDK直接轮询(推荐)

直接借助官方Python SDK提交任务并轮询状态,完成后触发后续操作,无需额外组件:

步骤示例

  1. 安装依赖:
pip install google-cloud-bigquery
  1. 编写代码实现完整流程:
from google.cloud import bigquery
import time

def load_avro_to_bq(gcs_uri, dataset_id, table_id):
    # 初始化BigQuery客户端
    client = bigquery.Client()
    
    # 配置Load任务参数
    job_config = bigquery.LoadJobConfig(
        source_format=bigquery.SourceFormat.AVRO,
        write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE  # 可根据需求改为WRITE_APPEND等
    )

    # 提交Load任务
    load_job = client.load_table_from_uri(
        gcs_uri,
        f"{dataset_id}.{table_id}",
        job_config=job_config
    )
    print(f"提交Load任务,Job ID: {load_job.job_id}")

    # 轮询任务状态,每10秒检查一次
    while not load_job.done():
        time.sleep(10)
        load_job.reload()  # 刷新最新状态
        print(f"任务当前状态: {load_job.state}")

    # 处理任务结果
    if load_job.error_result:
        raise Exception(f"Load任务失败: {load_job.error_result['message']}")
    print("Load任务执行完成")

    # 触发后续读取任务
    run_post_load_query(dataset_id, table_id)

def run_post_load_query(dataset_id, table_id):
    client = bigquery.Client()
    # 示例:读取表记录数
    query = f"SELECT COUNT(*) AS record_count FROM `{client.project}.{dataset_id}.{table_id}`"
    query_job = client.query(query)
    
    # 获取查询结果
    for row in query_job.result():
        print(f"目标表当前记录数: {row.record_count}")

# 执行示例
if __name__ == "__main__":
    load_avro_to_bq(
        gcs_uri="gs://your-bucket-name/path/to/*.avro",
        dataset_id="your_target_dataset",
        table_id="your_target_table"
    )

二、用BQ命令行+Shell脚本轮询

如果偏好命令行工具,可结合Shell脚本实现任务提交与状态监控:

#!/bin/bash

# 配置参数
GCS_URI="gs://your-bucket/path/*.avro"
DATASET_TABLE="your_dataset.your_table"

# 提交Load任务并提取Job ID
JOB_ID=$(bq load --source_format=AVRO "$DATASET_TABLE" "$GCS_URI" | grep -oP 'Job ID: \K\S+')
echo "提交Load任务,Job ID: $JOB_ID"

# 轮询任务状态
while true; do
    # 获取任务当前状态
    JOB_STATE=$(bq show --job=true "$JOB_ID" --format=json | jq -r '.status.state')
    
    case "$JOB_STATE" in
        "DONE")
            # 检查是否存在错误
            ERROR_MSG=$(bq show --job=true "$JOB_ID" --format=json | jq -r '.status.errorResult.message')
            if [ "$ERROR_MSG" != "null" ]; then
                echo "Load任务失败: $ERROR_MSG"
                exit 1
            fi
            echo "Load任务完成"
            # 执行后续读取操作,示例:查询表记录数
            bq query "SELECT COUNT(*) FROM $DATASET_TABLE"
            break
            ;;
        "FAILED")
            echo "Load任务执行失败"
            exit 1
            ;;
        *)
            echo "任务进行中,当前状态: $JOB_STATE"
            sleep 10
            ;;
    esac
done

三、结合调度工具(如Airflow)自定义传感器

如果用Airflow做任务调度,可自定义传感器监控Job状态,实现任务依赖:

from airflow.sensors.base import BaseSensorOperator
from google.cloud import bigquery
from airflow.utils.decorators import apply_defaults

class BigQueryLoadJobSensor(BaseSensorOperator):
    @apply_defaults
    def __init__(self, job_id, *args, **kwargs):
        super().__init__(*args, **kwargs)
        self.job_id = job_id

    def poke(self, context):
        client = bigquery.Client()
        job = client.get_job(self.job_id)
        if job.done():
            if job.error_result:
                raise Exception(f"Load任务失败: {job.error_result['message']}")
            return True
        # 任务未完成,继续等待
        return False

# 在DAG中定义任务依赖
from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator
from airflow.providers.google.cloud.operators.bigquery import BigQueryQueryOperator
from datetime import datetime

with DAG(
    dag_id="bq_avro_load_pipeline",
    start_date=datetime(2024, 1, 1),
    schedule_interval="@daily",
    catchup=False
) as dag:
    # 提交Load任务
    load_task = BigQueryInsertJobOperator(
        task_id="submit_avro_load",
        configuration={
            "load": {
                "sourceUris": ["gs://your-bucket/path/*.avro"],
                "destinationTable": {
                    "projectId": "your-gcp-project",
                    "datasetId": "your-dataset",
                    "tableId": "your-table"
                },
                "sourceFormat": "AVRO",
                "writeDisposition": "WRITE_TRUNCATE"
            }
        }
    )

    # 监控Load任务完成
    wait_for_load = BigQueryLoadJobSensor(
        task_id="wait_for_load_complete",
        job_id="{{ ti.xcom_pull(task_ids='submit_avro_load')['jobReference']['jobId'] }}"
    )

    # 后续读取任务
    post_load_query = BigQueryQueryOperator(
        task_id="post_load_count_query",
        sql="SELECT COUNT(*) FROM `your-gcp-project.your-dataset.your-table`"
    )

    # 设置任务依赖
    load_task >> wait_for_load >> post_load_query

内容的提问来源于stack exchange,提问作者salvob

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 23:10:45