如何在GCP Airflow中便捷监控BQ Load任务的Job ID状态?
实现BigQuery Load任务状态监控的简便方案
一、用BigQuery Python SDK直接轮询(推荐)
直接借助官方Python SDK提交任务并轮询状态,完成后触发后续操作,无需额外组件:
步骤示例
- 安装依赖:
pip install google-cloud-bigquery
- 编写代码实现完整流程:
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
相关产品推荐
相关产品推荐

