【发布时间】: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)
【问题讨论】: