如何在Airflow DAG中使用触发作业传入的JSON配置值
报错根因
dag_run是Airflow的运行时上下文变量,仅在DAG被实际触发执行的阶段才会生成。你当前拼接CLUSTER_NAME的代码位于DAG文件顶层,这段逻辑会被Airflow调度器在周期解析DAG文件时提前执行,此时上下文中不存在dag_run对象,因此抛出未定义错误。
修复方案
Airflow官方算子原生支持Jinja模板参数渲染,无需在DAG解析阶段提前拼接集群名称,直接将带模板语法的字符串传入cluster_name参数即可。任务实际运行时,Airflow会自动从运行上下文提取dag_run.conf中的配置值完成渲染,最终生成符合预期的集群名。
修正后的代码示例:
from airflow.utils.dates import days_ago from airflow import models # 替换为你实际使用的Dataproc集群创建算子导入路径 from airflow.providers.google.cloud.operators.dataproc import DataprocCreateClusterOperator CONN_ID = 'blah' PROJECT_ID = 'xyz' REGION = 'us-east4' # 定义集群名模板,运行时自动渲染 CLUSTER_NAME_TEMPLATE = "my-cluster-for-{{ dag_run.conf['x'] }}-{{ dag_run.conf['y'] }}-{{ dag_run.conf['z'] }}" with models.DAG('simple-python-dag', start_date=days_ago(1), schedule_interval=None) as dag: create_cluster_spark = DataprocCreateClusterOperator( task_id='create_cluster_spark', cluster_name=CLUSTER_NAME_TEMPLATE, region=REGION, gcp_conn_id=CONN_ID, # 需补充完整集群配置,否则调用Dataproc接口会失败 cluster_config={} )
注意事项
- 如果你使用的是自定义封装的
XyzCreateClusterOperator,需要将cluster_name加入算子类的template_fields元组,否则模板不会被渲染,示例:
class XyzCreateClusterOperator(BaseOperator): template_fields = ('cluster_name',) # 其余算子逻辑省略
- 如果需要在Python自定义函数中读取DAG运行配置,不要在顶层取值,要通过上下文注入的方式在任务执行阶段获取,示例:
def generate_cluster_name(**context): conf = context['dag_run'].conf return f"my-cluster-for-{conf['x']}-{conf['y']}-{conf['z']}"
Dataproc集群名称需符合GCP资源命名规则:仅支持小写字母、数字、连字符,不能以连字符开头或结尾,长度不超过55字符,不符合规则会导致集群创建失败。
内容的提问来源于stack exchange,提问作者Krishna Chaitanya V
相关产品推荐
相关产品推荐

