【问题标题】:How to return Object after evaluating the Templates in Airflow?评估气流中的模板后如何返回对象?
【发布时间】:2021-07-26 13:05:04
【问题描述】:

我们正在设计一个变量选择和参数设置器逻辑,当 DAG 被触发时需要对其进行评估。我们的 DAG 是在执行之前生成的。我们决定将我们的静态代码修改为自定义宏。


直到此时,在运算符定义之间定义了一个代码,因此当 DAG 由 DAG 生成器代码生成时,它正在运行。此代码无法处理用于选择正确 Airflow 变量的运行时参数。

for table_name in ast.literal_eval(Variable.get('PYTHON_LIST_OF_TABLES')):
    dag_id = "TableLoader_" + str(table_name)
    default_dag_args={...}
    schedule = None
    globals()[dag_id] = create_dag(dag_id, schedule, default_dag_args)
def create_dag(dag_id, schedule, default_dag_args):

    with DAG(
        default_args=default_dag_args,
        dag_id=dag_id,
        schedule_interval=schedule,
        user_defined_macros={ "load_type_handler": load_type_handler }
    ) as dag:

        # static python code which sets pipeline_rt_args for all generated DAGs the same way
        # this static code could set only one type (INITIAL or INCREMENTAL)
        # but we want to decide during the execution now

        # Operator Definitions
        OP_K = CloudDataFusionStartPipelineOperator(
            task_id='PREFIX_'+str(table),

            # ---> Can't handle runtime parameters <---
            runtime_args=pipeline_rt_args,                             
            # ...
        )

        OP_1 >> OP_2 >> ... >> OP_K >> ... >> OP_N
        return dag

现在我们想要在从 UI 或 REST API 触发 DAG 时传递 load_type(例如:INITIALINCREMENTAL),因此我们需要修改这个旧的(静态)行为(它只处理一种情况,但不是两种情况)来获取正确的气流变量并为我们的CloudDataFusionStartPipelineOperator创建正确的对象:

例如:

{"load_type":"INCREMENTAL"}
# or
{"load_type":"INITIAL"}

但如果我们这样做:


def create_dag(dag_id, schedule, default_dag_args):

    def extend_runtime_args(prefix, param, field_name, date, job_id):

        # reading the Trigger-time parameter
        load_type = param.conf["load_type"]
        
        # getting the proper Airflow Variable (depending on current load type)
        result = eval(Variable.get(prefix+'_'+load_type+'_'+dag_id))[field_name]

        # setting 'job_id', 'dateValue', 'date', 'GCS_Input_Path' for CloudDataFusionStartPipelineOperator
        # ...

        return rt_args


    with DAG( #...
        user_defined_macros={
            "extend_runtime_args": extend_runtime_args
        }) as dag:
        # removed static code (which executes only in generation time)

        # Operator Definitions
        OP_K = CloudDataFusionStartPipelineOperator(
            task_id='PREFIX_'+str(table),

            # ---> handles runtime arguments with custom macro <---
            runtime_args="""{{ extend_runtime_args('PREFIX', dag_run, 'runtime_args', macros.ds_format(yesterday_ds_nodash,"%Y%m%d","%Y_%m_%d"), ti.job_id) }}""",
            # ...
        )

        OP_1 >> OP_2 >> ... >> OP_K >> ... >> OP_N
        return dag

注意:

这里我们需要的是自定义逻辑的“未来”评估(不在 DAG 生成时评估),它将返回一个对象,这就是我们在这里使用模板的原因。

我们会遇到以下情况:

  • 在自定义宏函数extend_runtime_args内部,返回类型是一个对象
  • 评估 Jinja 模板后,返回类型更改为字符串
  • CloudDataFusionStartPipelineOperator 失败,因为 runtime_args 属性是字符串而不是对象

问题:

  • 我们如何在评估 Jinja 模板后返回一个对象(并在“未来”这样做)?
    • 我们能以某种方式转换字符串吗?
  • 如何保证这里的逻辑是在 DAG 执行后执行,而不是在 DAG 生成后执行?
  • Jinja 模板/自定义宏在这里处理触发时间参数是好还是坏?

【问题讨论】:

标签: python jinja2 airflow google-cloud-data-fusion


【解决方案1】:

评估 Jinja 模板后,我们如何返回一个对象 (并在“未来”这样做)?

您可以创建自己的从CloudDataFusionStartPipelineOperator 派生的自定义运算符,并使其接受字符串并将其转换为CloudDataFusionStartPipelineOperator 所需的对象并使用此新运算符。 “runtime_args”是一个字典,所以我相信它应该像 json.loads() 一样容易找回它。

我们可以以某种方式转换字符串吗?

是的。只需上面的 json.loads() 代码即可。此外,如果您在 runtime_args 中只有几个参数会更改,那么拥有多个宏并直接在字典中的多个 JINJA 字符串中返回更改后的值会更容易。比如:

runtime_args = {
   'PREFIX' = "{{ dag_run }}",
   'date' = "{{ macros.ds_format(....) }}",
}

当存在模板字段时,Airflow 会递归处理字典或列表等基本结构,因此您可以保留对象结构,并将 jinja 宏用作值(实际上您也可以将 jinja 宏用作键等)。

我们如何保证这里的逻辑会在 DAG 完成后执行 在生成 DAG 之后执行而不是立即执行?

仅在执行任务时评估 JINJA 模板。所以你在这里很好。

Jinja 模板/自定义宏在这里是好还是坏模式? 处理触发时间参数?

非常好的模式。这就是他们的目的。

【讨论】:

  • 谢谢!这是一个非常好的解决方案......为了记录,我们已经以不同的方式解决了这个问题:我们正在逐个字段地组装runtime_args 对象。字段接受字符串,如果单独设置则不需要返回对象。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-03-24
  • 2019-07-29
  • 2020-01-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多