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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 23:07:28