【问题标题】:Airflow: how to use trigger parameters in functions气流:如何在函数中使用触发参数
【发布时间】:2021-09-07 05:03:03
【问题描述】:

我们将 Airflow 的 KubernetesPodOperator 用于我们的数据管道。我们要添加的是通过 UI 传递参数的选项。

我们目前使用它的方式是,我们有不同的 yaml 文件来存储操作符的参数,我们不是直接调用操作符,而是调用一个函数来进行一些准备并返回操作符,如下所示:

def prep_kubernetes_pod_operator(yaml):

    # ... read yaml and extract params

    return KubernetesPodOperator(params)

with DAG(...):
    
    task1 = prep_kubernetes_pod_operator(yaml)

对我们来说,这很好用,我们可以保持我们的 dag 文件非常轻量级,但是现在我们想添加可以通过 UI 添加一些额外参数的功能。我知道可以通过kwargs['dag_run'].conf 访问触发参数,但我没有成功将它们拉入 Python 函数。

我尝试的另一件事是创建一个自定义运算符,因为它可以识别 args,但我无法在执行部分调用 KubernetesPodOperator(而且我想在运算符中调用运算符无论如何都不是正确的解决方案)。

更新:

按照 NicoE 的建议,我开始改为扩展 KubernetesPodOperator

我现在遇到的错误是,当我解析 yaml 并在之后分配参数时,父参数变为元组并引发类型错误。

dag:

task = NewKPO(
    task_id="task1",
    yaml_path=yaml_path)

运营商:

class NewKPO(KubernetesPodOperator):
   @apply_defaults
   def __init__(
           self,
           yaml_path: str,
           name: str = "default",
           *args,
           **kwargs) -> None:
       self.yaml_path = yaml_path
       self.name = name
       super(NewKPO, self).__init__(
           name=name, # DAG is not parsed without this line - 'key has to be string'
           *args,
           **kwargs)

   def execute(self, context):
       # parsing yaml and adding context["dag_run"].conf (...)
       self.name = yaml.name
       self.image = yaml.image
       self.secrets = yaml.secrets
       #(...) if i run a type(self.secrets) here I will get tuple
       return super(NewKPO, self).execute(context)

【问题讨论】:

    标签: python airflow


    【解决方案1】:

    您可以使用params,它是一个可以在 DAG 级别参数定义的字典,并且在每个任务中都可以访问。适用于从BaseOperator 派生的每个运算符,也可以从 UI 中设置。

    以下示例显示了如何将其与不同的运算符一起使用。 params 可以在 default_args dict 中定义,也可以作为 DAG 对象的 arg。

    default_args = {
        "owner": "airflow",
        'params': {
            "param1": "first_param",
            "param2": "second_param"
        }
    }
    
    dag = DAG(
        dag_id="example_dag_params",
        default_args=default_args,
        start_date=days_ago(1),
        schedule_interval="@once",
        tags=['example_dags'],
        catchup=False
    )
    
    

    从 UI 触发此 DAG 时,您可以添加一个额外的参数:

    可以在模板化字段中访问参数,如 BashOperator 案例:

    with dag:
    
        bash_task = BashOperator(
            task_id='bash_task',
            bash_command='echo bash_task: {{ params.param1 }}')
    
    

    bash_task 日志输出:

    {bash.py:158} INFO - Running command: echo bash_task: first_param
    {bash.py:169} INFO - Output:
    {bash.py:173} INFO - bash_task: first_param
    {bash.py:177} INFO - Command exited with return code 0
    

    参数可以在执行上下文中访问,例如python_callable:

    
        def _print_params(**kwargs):
            print(f"Task_id: {kwargs['ti'].task_id}")
            for k, v in kwargs['params'].items():
                print(f"{k}:{v}")
    
        python_task = PythonOperator(
            task_id='python_task',
            python_callable=_print_params,
        )
    

    输出:

    {logging_mixin.py:104} INFO - Task_id: python_task
    {logging_mixin.py:104} INFO - param1:first_param
    {logging_mixin.py:104} INFO - param2:second_param
    {logging_mixin.py:104} INFO - param3:param_from_the_UI
    

    您还可以在任务级别定义中添加参数:

        python_task_2 = PythonOperator(
            task_id='python_task_2',
            python_callable=_print_params,
            params={'param4': 'param defined at task level'}
        )
    

    输出:

    {logging_mixin.py:104} INFO - Task_id: python_task_2
    {logging_mixin.py:104} INFO - param1:first_param
    {logging_mixin.py:104} INFO - param2:second_param
    {logging_mixin.py:104} INFO - param4:param defined at task level
    {logging_mixin.py:104} INFO - param3:param_from_the_UI
    

    按照示例,您可以定义一个继承自 BaseOperator 的自定义 Operator:

    class CustomDummyOperator(BaseOperator):
    
        @apply_defaults
        def __init__(self, custom_arg: str = 'default', *args, **kwargs) -> None:
            self.custom_arg = custom_arg
            super(CustomDummyOperator, self).__init__(*args, **kwargs)
    
        def execute(self, context):
            print(f"Task_id: {self.task_id}")
            for k, v in context['params'].items():
                print(f"{k}:{v}")
    

    一个示例任务是:

        custom_op_task = CustomDummyOperator(
            task_id='custom_operator_task'
        )
    

    输出:

    {logging_mixin.py:104} INFO - Task_id: custom_operator_task
    {logging_mixin.py:104} INFO - custom_arg: default
    {logging_mixin.py:104} INFO - param1:first_param
    {logging_mixin.py:104} INFO - param2:second_param
    {logging_mixin.py:104} INFO - param3:param_from_the_UI
    

    进口:

    from airflow import DAG
    from airflow.models.baseoperator import chain
    from airflow.models import BaseOperator
    from airflow.operators.python import PythonOperator
    from airflow.operators.bash import BashOperator
    from airflow.utils.dates import days_ago
    from airflow.utils.decorators import apply_defaults
    

    希望对你有用!

    【讨论】:

    • 嘿 Nico,我非常感谢您的彻底回答。虽然你所说的完全有道理,但这并不是我想要的。我会使用 UI 添加额外的参数,但我希望我编写的 Python 函数 (prep_kubernetes_pod_operator) 作为示例来获取它们。所以这不会是 PythonOperator 的可调用对象,因为最终我将运行 KubernetesPodOperator。我希望这次我能够正确解释它。这种用法会有解决方案吗?非常感谢。
    • @a54i 好的,现在我明白了!对不起。我认为没有一种方法可以像您提供的示例中那样从“任意”函数访问paramsconf。所以我想到了两个选项,使用KubernetesPodOperator 中的template fields,它们是:'image', 'cmds', 'arguments', 'env_vars', 'labels', 'config_file', 'pod_template_file' with jinja 模板,与上面的 BashOperator 示例中所示的方式相同。
    • 另一种方法,如果您需要访问这些参数,可能会处理它们,并将它们作为参数传递给KubernetesPodOperator,但在template_fields之外,那么您可以考虑创建您的扩展 KubernetesPodOperator 的自定义操作符。这将允许您做几乎任何您需要的事情,并通过实例化这些自定义运算符直接创建任务,删除函数。让我知道这是否对您有用!
    • 谢谢 Nico,是的,我相信最好的解决方案是创建一个自定义运算符。那么我是否只需创建一个从 KubernetesPodOperator 继承的新运算符,然后我可以更改一些参数,如 self.image 等,然后父运算符会执行吗?抱歉,我对这些不太熟悉。
    • 问题并不特定于秘密,我分配的所有参数都变成了元组。我会按照你的建议提出一个新问题,再次感谢 Nico!
    猜你喜欢
    • 2020-05-14
    • 2021-05-14
    • 2021-05-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多