如何将Airflow运行日期(ds)作为变量传递至GCP Dataform
解决方案:从Airflow传递运行日期到Dataform管道
问题根源
你之前的错误是因为DataformCreateCompilationResultOperator的compilation_result对象并不包含vars字段,变量传递应该通过**compilation_config参数**实现,这是Dataform编译配置的专属字段。
步骤1:确认Dataform项目配置
保持你的dataform.json中的默认vars配置(作为 fallback 值,避免无传参时出错):
{ "defaultSchema": "dataform", "assertionSchema": "dataform_assertions", "warehouse": "bigquery", "defaultDatabase": "prod-dna-pipelines", "defaultLocation": "EU", "vars" : { "date" : "2023-05-01" } }
步骤2:修改Airflow Operator代码
将变量传递逻辑移到compilation_config参数中,替换你原来的错误写法:
create_compilation_result = DataformCreateCompilationResultOperator( task_id="create_compilation_result", project_id=PROJECT_ID, region=REGION, repository_id=REPOSITORY_ID, compilation_result={ "git_commitish": GIT_COMMITISH, "workspace": ( f"projects/{PROJECT_ID}/locations/{REGION}/repositories/{REPOSITORY_ID}/" f"workspaces/{WORKSPACE_ID}" ) }, # 新增compilation_config参数传递变量 compilation_config={ "vars": { "date": "{{ ds }}" } } )
{{ ds }}是Airflow的模板变量,会自动替换为DAG的运行日期(格式YYYY-MM-DD)- 未来需要传递其他变量时,直接在
vars字典中添加键值对即可,比如"environment": "{{ var.value.dataform_env }}"
步骤3:在Dataform SQLX文件中使用变量
在你的Dataform SQL脚本中,通过${vars.date}引用传递过来的日期变量,替换原来的CURRENT_DATE:
config { type: "table" } SELECT * FROM `prod-dna-pipelines.your_dataset.your_source_table` WHERE date = ${vars.date}
验证重跑功能
当你在Airflow中触发历史日期的重跑任务时,{{ ds }}会被替换为对应的历史日期,Dataform会使用该日期查询源表,完美解决重跑需求。
内容的提问来源于stack exchange,提问作者kfondingle
相关产品推荐
相关产品推荐

