【问题标题】:airflow - dynamically change namespace in KubernetesPodOperator气流 - 在 KubernetesPodOperator 中动态更改命名空间
【发布时间】:2021-03-17 07:58:33
【问题描述】:

我在 kubernetes 上运行气流 1.10.13 并试图找到一种方法来动态更改我运行任务的命名空间。 我尝试使用模板并从 dag_run.conf json 插入参数,但模板仅在“cmds”中呈现,而不在命名空间等其他任务字段中呈现。

我很想找到一个解决方案(使用模板或任何其他方式)来更改命名空间。

default_args = {
    'owner': 'airflow',
    'start_date': days_ago(1)
}

with DAG('test_ns', default_args=default_args, schedule_interval='@once') as dag:
    ns = """ {{ dag_run.conf.ns }} """
    example_task= KubernetesPodOperator(namespace=ns,
                                         image='python:3.6',
                                         cmds=["/bin/sh", "-c", "echo {{ns}}"],
                                         arguments=[],
                                         task_id='example_task',
                                         name='example_task',
                                         get_logs=True,
                                         is_delete_operator_pod=True,
                                         provide_context=True
                                         )

【问题讨论】:

    标签: templates macros airflow


    【解决方案1】:

    KubernetesPodOperator 中templated fields 的列表没有namespace

    您可以创建自己的运算符,其行为与将namespace 添加到template_fields 的行为相同:

    class MyKubernetesPodOperator(KubernetesPodOperator):
             template_fields = KubernetesPodOperator.template_fields +('namespace',)
        
    

    然后你可以使用你的代码:

    with DAG('test_ns', default_args=default_args, schedule_interval='@once') as dag:
        ns = """ {{ dag_run.conf.ns }} """
        example_task= MyKubernetesPodOperator(namespace=ns,...)
    

    编辑: 您的代码的完整示例:

    from datetime import datetime
    from airflow import DAG
    from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator
    
    
    class MyKubernetesPodOperator(KubernetesPodOperator):
        template_fields = KubernetesPodOperator.template_fields + ('namespace',)
    
    
    default_args = {
        'owner': 'elad',
        'start_date': datetime(2019, 11, 1),
    
    }
    
    with DAG(dag_id='stackoverflow',
             default_args=default_args,
             schedule_interval=None
             ) as dag:
        ns = """ {{ dag_run.conf.ns }} """
        example_task = MyKubernetesPodOperator(namespace=ns,
                                               image='python:3.6',
                                               cmds=["/bin/sh", "-c", "echo {{ns}}"],
                                               arguments=[],
                                               task_id='example_task',
                                               name='example_task',
                                               get_logs=True,
                                               is_delete_operator_pod=True,
                                               )
    

    conf触发dag:

    {"ns":"mynamespace"}
    

    【讨论】:

    • 感谢 Elad,我明白了这个想法,但在实施时遇到了麻烦。我正在尝试找到创建自定义运算符的正确语法。 class MyOperator(KubernetesPodOperator): @apply_defaults def __init__(self, *args, **kwargs): super(MyOperator, self).__init__(*args, **kwargs) template_fields = ('namespace') + KubernetesPodOperator.template_fields 我在教程中看到的这个版本为 init 返回了一个错误:__init__() got an unexpected keyword argument 'limit_memory
    • 为什么加init?您是否尝试过我答案中的代码?它应该可以正常工作。
    • 如果按原样使用,则会出现无效的语法错误。我使用版本 1.10.13。使用类是我尝试实现一个新的运算符,就像我在这里看到的那样:airflow.apache.org/docs/apache-airflow/stable/howto/…
    • 从哪里导入 KubernetesPodOperator?贡献者还是提供者?
    • 哦,它应该是('namespace',),所以它将是一个元组而不是一个字符串。我很快就会用工作示例更新我的答案。
    猜你喜欢
    • 1970-01-01
    • 2019-03-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-10-14
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多