如何为BigQuery创建数据感知型DAG?能否将BigQuery表设为Airflow数据集?
用BigQuery表作为Airflow Dataset实现数据感知触发
当然可以将BigQuery表设置为Airflow的Dataset,以此实现当表数据变更时触发依赖的DAG或任务,具体实现方式如下:
1. 定义BigQuery Dataset对象
使用Airflow的Dataset类,通过BigQuery表的URI来声明数据集,格式为bigquery://<GCP项目ID>/<BigQuery数据集ID>/<目标表ID>:
from airflow import Dataset # 替换为你的实际GCP项目、BQ数据集和表名 bq_target_dataset = Dataset("bigquery://my-gcp-project/my-bq-dataset/my-monitored-table")
2. 标记更新BigQuery表的任务
在负责写入/更新目标BigQuery表的Airflow任务中,通过outlets参数标记该任务会更新上述Dataset。以BigQueryInsertJobOperator为例:
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator update_bq_table = BigQueryInsertJobOperator( task_id="update_bq_table", configuration={ "query": { "query": "INSERT INTO my-gcp-project.my-bq-dataset.my-monitored-table SELECT * FROM staging_table", "useLegacySql": False, } }, outlets=[bq_target_dataset], # 标记该任务更新了目标Dataset )
3. 配置依赖DAG的触发规则
在需要被触发的依赖DAG中,将schedule参数设置为定义好的BigQuery Dataset,这样当Dataset被更新时,该DAG会自动触发:
from airflow import DAG from datetime import datetime with DAG( dag_id="bq_dependent_workflow", start_date=datetime(2024, 1, 1), schedule=[bq_target_dataset], # 监听Dataset变更触发DAG catchup=False, ) as dag: # 这里添加你的任务逻辑 pass
4. 外部系统更新表的处理
如果BigQuery表是被Airflow之外的系统更新的(比如ETL工具、手动SQL操作),需要手动触发Airflow的Dataset事件:
- 可以通过Airflow API调用触发:
curl -X POST "http://<airflow-webserver-url>/api/v1/datasets/bigquery%3A%2F%2Fmy-gcp-project%2Fmy-bq-dataset%2Fmy-monitored-table/events" \ -H "Authorization: Bearer <airflow-api-token>"
- 或者用Python SDK在监听脚本中调用(比如搭配Cloud Function监听BigQuery的变更日志):
from airflow.api.client.local_client import Client client = Client(None, None) client.trigger_dataset_event( dataset_uri="bigquery://my-gcp-project/my-bq-dataset/my-monitored-table" )
注意事项
- 确保你的Airflow Google Cloud Provider版本不低于10.0.0,该版本才正式支持BigQuery作为Dataset的数据源。
- 若使用工作负载身份验证,需确保Airflow服务账号拥有BigQuery的读取权限以及Airflow API的调用权限(如果涉及外部触发)。
内容的提问来源于stack exchange,提问作者Amarjeet Kushwaha
相关产品推荐
相关产品推荐

