如何在GCP Composer的Airflow DAG运行后向BigQuery写入状态数据
实现方案
1. 定义BigQuery日志表结构
你可以手动在BigQuery中创建日志表,也可以让DAG自动完成表的初始化。表结构如下:
CREATE TABLE IF NOT EXISTS `你的GCP项目ID.目标数据集ID.dag_run_logs` ( dag_run_start_time TIMESTAMP NOT NULL, dag_id STRING NOT NULL, run_state INT64 NOT NULL );
字段说明:
dag_run_start_time:DAG启动时间dag_id:DAG的唯一标识(方便后续区分不同DAG的日志)run_state:运行状态,0代表成功,1代表失败
2. Airflow DAG代码实现
在你的DAG中添加两个核心任务:表初始化任务、运行状态写入任务。利用Airflow上下文变量获取DAG运行元数据,通过BigQueryInsertJobOperator完成数据追加。
from airflow import DAG from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator, BigQueryCreateEmptyTableOperator from airflow.utils.dates import days_ago from airflow.utils.trigger_rule import TriggerRule # 替换为你的GCP配置 PROJECT_ID = "你的GCP项目ID" DATASET_ID = "目标数据集ID" TABLE_ID = "dag_run_logs" default_args = { 'owner': 'airflow', 'start_date': days_ago(1), } with DAG( '你的DAG名称', default_args=default_args, description='每5分钟运行并记录状态到BigQuery', schedule_interval='*/5 * * * *', catchup=False, ) as dag: # 任务1:确保日志表存在,不存在则自动创建 init_log_table = BigQueryCreateEmptyTableOperator( task_id='init_log_table', project_id=PROJECT_ID, dataset_id=DATASET_ID, table_id=TABLE_ID, schema_fields=[ {'name': 'dag_run_start_time', 'type': 'TIMESTAMP', 'mode': 'REQUIRED'}, {'name': 'dag_id', 'type': 'STRING', 'mode': 'REQUIRED'}, {'name': 'run_state', 'type': 'INT64', 'mode': 'REQUIRED'} ], exists_ok=True, ) # 生成插入日志的SQL语句 def generate_insert_sql(**context): dag_run = context['dag_run'] start_time = dag_run.start_date.isoformat() dag_id = dag_run.dag_id run_state = 0 if dag_run.state == 'success' else 1 return f""" INSERT INTO `{PROJECT_ID}.{DATASET_ID}.{TABLE_ID}` (dag_run_start_time, dag_id, run_state) VALUES (TIMESTAMP('{start_time}'), '{dag_id}', {run_state}) """ # 任务2:写入DAG运行状态到BigQuery log_run_status = BigQueryInsertJobOperator( task_id='log_run_status', project_id=PROJECT_ID, configuration={ "query": { "query": "{{ generate_insert_sql() }}", "useLegacySql": False, } }, trigger_rule=TriggerRule.ALL_DONE, # 无论DAG成功/失败都执行 ) # 设置任务依赖 init_log_table >> log_run_status
3. 关键配置说明
- TriggerRule.ALL_DONE:确保不管DAG中其他任务(如果有)执行结果如何,状态日志都会被写入
- 上下文变量:通过
context['dag_run']直接获取DAG运行的启动时间、状态等元数据,无需额外处理 - 权限配置:确保Composer的服务账号拥有目标BigQuery数据集的
dataEditor权限,或者以下精细权限:bigquery.tables.create(自动建表需要)bigquery.jobs.create(执行插入查询需要)bigquery.tables.updateData(写入数据需要)
内容的提问来源于stack exchange,提问作者Dipak Wani
相关产品推荐
相关产品推荐

