【问题标题】:Condition based on BranchPythonOperator基于 BranchPythonOperator 的条件
【发布时间】:2022-10-15 01:06:46
【问题描述】:

我有一些表单的自定义运算符

class DataPreparationOperator(BaseOperator):

template_fields = ['file']

def __init__(self, file, **kwargs):
    super().__init__(**kwargs)
    self.file = file


...

def execute(self, context):
    filename = f'prepared_data_{str(time()*1000).replace(".","_")}.json'
    DataDownloader(filename, self.filters()).dataframe_downloader()
    
    return filename

(...)

class DataPreparationOperatorArrivals(DataPreparationOperator):

template_fields = ['file']

def __init__(self, file, **kwargs):
    super(DataPreparationOperatorArrivals, self).__init__(file=file, **kwargs)

...
def execute(self, context):
    filename = f'prepared_data_{str(time()*1000).replace(".","_")}.json'
    DataDownloader(filename, self.charge_change()).dataframe_downloader()
    
    return filename

(...)

运算符基于 BranchPythonOperator 执行,在我的 dag 中如下所示

def choose_data_preparation_operator(**kwargs):
    if float(kwargs.get("arrival_factor")) != 1.0:
        return ['data_preparation_arrivals_change', 'parameters_constructor']
    else:
        return ['data_preparation_normal', 'parameters_constructor']

opr_data_preparation_path = BranchPythonOperator(
    provide_context=True,
    task_id='choose_data_preparation_path',
    python_callable=choose_data_preparation_operator,
    op_kwargs = {'arrival_factor': '{{ dag_run.conf["arrival_factor"] }}'}
    )

opr_data_prep = DataPreparationOperator(
    task_id ='data_preparation_normal',
    file = 'data.json'
    )

opr_data_prep_arr = DataPreparationOperatorArrivals(
    task_id ='data_preparation_arrivals_change',
    file = 'data.json'
    )

如您所见,两个运算符都返回一个文件名,现在我想使用另一个自定义运算符并使用各自的文件名在另一个步骤中调用此文件,例如

opr_parameters_constructor = ParametersConstructor(
    task_id ='parameters_constructor',
    file = '{{ ti.xcom_pull(task_ids="CHOOSE_THE_CORRECT_TASK_ID") }}',
    initial_time = '{{ dag_run.conf.get("initial_time") }}',
    final_time = '{{ dag_run.conf.get("final_time") }}',
    )

我的问题是,我怎样才能把正确的task_id在 BranchPythonOperator? 中选择,即最后一段代码中的 CHOOSE_THE_CORRECT_TASK_ID 变量。

非常感谢您的帮助:D

【问题讨论】:

    标签: python airflow directed-acyclic-graphs


    【解决方案1】:

    我认为嵌套表达式可能对你有用。在内部表达式中使用ti.xcom_pull()BranchPythonOperatortask_id 并索引到数组以获取返回文件名的任务的task_id

    opr_parameters_constructor = ParametersConstructor(
        task_id ='parameters_constructor',
        file = '{{{{ ti.xcom_pull(task_ids="{{ ti.xcom_pull(task_ids="choose_data_preparation_path")[0] }}" }}}}',
        initial_time = '{{ dag_run.conf.get("initial_time") }}',
        final_time = '{{ dag_run.conf.get("final_time") }}',
    )
    

    让我知道这是否适合您。

    【讨论】:

    • 嗨,感谢您的帮助,我有一个问题,当我触发 dag 时,显示以下错误 '''jinja2.exceptions.TemplateSyntaxError: expected token ':', got '}''''
    • 嗨 0xTochi,我编辑了解决方案,它缺少一些引号。你能再测试一下,看看它是否适合你吗?谢谢。
    【解决方案2】:

    我最近遇到了类似的挑战。

    我在第一个任务中创建了 xcom_push 操作,以指示每个单独任务的文件应该放在哪里。然后,以下步骤将使用 xcom_pull 操作来获取文件位置。最后,我扩展了操作符pre_execute 函数,再次用 xcom_pulls 覆盖了一些变量。

    但是,我相信在您的情况下,创建以分离同一运算符的实例要简单得多。这样您就不必根据运行的任务进行复杂的匹配和变量更改。相反,您只需编写两个运算符,每个运算符都匹配两个任务之一。

    【讨论】:

      猜你喜欢
      • 2012-06-29
      • 2023-03-11
      • 1970-01-01
      • 2014-09-20
      • 1970-01-01
      • 2018-10-27
      • 2013-11-12
      • 2016-01-15
      • 2012-09-11
      相关资源
      最近更新 更多