【问题标题】:KubernetesPodOperator - run entire command pulled from XCOMKubernetesPodOperator - 运行从 XCOM 提取的整个命令
【发布时间】:2021-10-11 05:04:35
【问题描述】:

我正在开发 Airflow 1.10。

我在 KubernetesPodOperator 上运行命令时遇到问题,在 DAG 运行时评估整个命令。

我在 DAG 运行时生成命令,因为一些命令的参数取决于用户传递的参数。

正如我从文档中看到的 KubernetesPodOperator 需要字符串列表或 jinja 模板列表:

    :param arguments: arguments of the entrypoint. (templated)
        The docker image's CMD is used if this is not provided.

我有 PythonOperator,它生成命令并将其推送到 XCOM 和 KubernetesPodOperator 在参数中我传递 PythonOperator 生成的命令。

from airflow.operators.python_operator import PythonOperator
from airflow.contrib.operators.kubernetes_pod_operator import KubernetesPodOperator


def command_maker():
    import random # random is to illustrate that we don't know arguments value before runtime
    return f"my_command {random.randint(1, 10)} --option {random.randint(1, 4)}"

def create_tasks(dag):
    first = PythonOperator(
        task_id="generate_command",
        python_callable=command_maker,
        provide_context=True,
        dag=dag,
    )
    second = KubernetesPodOperator(
        namespace='some_namespace',
        image='some_image',
        name='execute_command',
        dag=dag,
        arguments=[f'{{ ti.xcom_pull(dag_id="{dag.dag_id}", task_ids="generate_command", key="return_value")}}']
    )
    second.set_upstream(first)

不幸的是 KubernetesPodOperator 没有正确运行这个命令,因为他试图运行这样的东西:

[my_command 4 --option 2]

有没有办法在 KubernetesPodOperator 运行时评估这个列表 或者我是否强制将所有运行时参数推送到单独的 XCOM 中? 我想避免这样的解决方案,因为它需要在我的项目中进行大量更改。

         arguments=[
            "my_command",
            f'{{ ti.xcom_pull(dag_id="{dag.dag_id}", task_ids="generate_command", key="first_argument")}}',
            "--option",
            f'{{ ti.xcom_pull(dag_id="{dag.dag_id}", task_ids="generate_command", key="second_argument")}}',
         ]

【问题讨论】:

    标签: python kubernetes airflow


    【解决方案1】:

    问题是JINJA模板默认返回模板为字符串。

    然而,在最近的 Airflow 中(从 Airlfow 2.1.0 开始),您可以将模板渲染为原生 python 对象:

    https://airflow.apache.org/docs/apache-airflow/stable/concepts/operators.html#rendering-fields-as-native-python-objects

    通过在创建 DAG 时使用 render_template_as_native_obj=True 参数。

    然后,您需要以 python 的literal_eval 能够将其转换为 python 对象的方式格式化输出。在您的情况下,您必须使输出类似于:

    [ 'my_command', '4', '--option', '2' ]

    请注意,此参数将为您的所有模板返回原生对象,因此如果它们返回一些 literal_eval understands 的值 - 它们也将被转换为原生类型(并且您可能会产生一些意想不到的副作用。

    【讨论】:

    • 在 Airflow 1.10 中是否有任何解决方法来实现这一点?
    【解决方案2】:

    要获得适用于 Airflow 1.10 的有效解决方案,我必须使用 BaseOperator.pre_execute 挂钩。

    from airflow.contrib.operators.kubernetes_pod_operator import KubernetesPodOperator
    from airflow.lineage import prepare_lineage
    
    class UnpackCommandKubernetesPodOperator(KubernetesPodOperator):
        @prepare_lineage
        def pre_execute(self, context):
            self.arguments = self.arguments[0].split(" ")
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-04-27
      • 2020-05-02
      • 1970-01-01
      • 2023-01-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-05-01
      相关资源
      最近更新 更多