【问题标题】:Task after BranchPythonOperator Task getting skippedBranchPythonOperator 任务被跳过后的任务
【发布时间】:2021-01-08 18:16:06
【问题描述】:

我创建了一个 BranchPythonOperator,它根据条件调用 2 个任务,例如:

typicon_check_table = BranchPythonOperator(
    task_id='typicon_check_table',
    python_callable=CheckTable(),
    provide_context=True,
    dag=typicon_task_dag)

typicon_create_table = PythonOperator(
    task_id='typicon_create_table',
    python_callable=CreateTable(),
    provide_context=True,
    dag=typicon_task_dag)

typicon_load_data = PythonOperator(
    task_id='typicon_load_data',
    python_callable=LoadData(),
    provide_context=True,
    dag=typicon_task_dag)

typicon_check_table.set_downstream([typicon_load_data, typicon_create_table])
typicon_create_table.set_downstream(typicon_load_data)

这是CheckTable 可调用类:

class CheckTable:
    """
    DAG task to check if table exists or not.
    """

    def __call__(self, **kwargs) -> None:
        pg_hook = PostgresHook(postgres_conn_id="postgres_docker")
        query = "SELECT EXISTS ( \
            SELECT 1 FROM information_schema.tables \
            WHERE table_schema = 'public' \
            AND table_name = 'users');"

        table_exists = pg_hook.get_records(query)[0][0]
        if table_exists:
            return "typicon_load_data"
        return "typicon_create_table"

问题是运行typicon_check_table 任务时这两个任务都被跳过了。

如何解决这个问题?

【问题讨论】:

    标签: airflow


    【解决方案1】:

    我已经解决了相同的场景,它适用于以下代码

    BranchPythonOperator(task_id='slot_population_on_is_y_or_n', python_callable=DAGConditionalValidation('Y'),
                             trigger_rule='one_success')
    slot_population_on_is_y = DummyOperator(task_id='slot_population_on_is_y')
    slot_population_on_is_n = DummyOperator(task_id='slot_population_on_is_n')
    slot_population_on_is_y_or_n >> [slot_population_on_is_y, slot_population_on_is_n]
    
    
    class DAGConditionalValidation:
    
        def __init__(self, conditional_param_key):
            self.conditional_param_key = conditional_param_key
    
    
        def __call__(self, **kwargs):
            if (conditional_param_key == 'Y'):
                return slot_population_on_is_y
            return slot_population_on_is_n
    

    您的所有代码看起来都很好,但是您缺少触发规则,请将触发规则设置为trigger_rule='one_success'。
    这应该也适合你。

    【讨论】:

    • 没错,我是新贡献者,所以需要学习发帖技巧:)
    【解决方案2】:

    任务typicon_load_data 有typicon_create_table 作为父任务,默认trigger_rule 是all_success,所以我对这种行为并不感到惊讶。

    这里有两种可能的情况:

    1. CheckTable() 返回typicon_load_data,然后typicon_create_table 被跳过,但typicon_load_data 下游也被跳过。
    2. CheckTable() 返回 typicon_create_table,它被执行并触发 typicon_load_data,因为它是被排除的分支,所以被跳过。

    我假设您的屏幕截图来自案例 1。?

    【讨论】:

      【解决方案3】:

      将 trigger_rule="all_done" 规则添加到 typicon_check_table 如下

      typicon_check_table = BranchPythonOperator(
          task_id='typicon_check_table',
          python_callable=CheckTable(),
          provide_context=True,
          trigger_rule="all_done",
          dag=typicon_task_dag)
      

      【讨论】:

        猜你喜欢
        • 2019-08-05
        • 2022-09-28
        • 1970-01-01
        • 1970-01-01
        • 2022-08-03
        • 1970-01-01
        • 1970-01-01
        • 2022-12-01
        • 1970-01-01
        相关资源
        最近更新 更多