【问题标题】:How does Airflow's BranchPythonOperator work?Airflow 的 BranchPythonOperator 是如何工作的?
【发布时间】:2017-11-15 12:26:53
【问题描述】:

我很难理解 Airflow 中的 BranchPythonOperator 是如何工作的。我知道它主要用于分支,但是文档对传递给任务的内容以及我需要从上游任务传递/期望的内容感到困惑。

鉴于文档on this page 中的简单示例,上游任务run_this_first 和分支的2 个下游任务的源代码是什么样的? Airflow 究竟是如何知道运行 branch_a 而不是 branch_b 的?上游任务的输出在哪里得到注意/读取?

【问题讨论】:

标签: airflow apache-airflow


【解决方案1】:

您的 BranchPythonOperator 是使用 python_callable 创建的,这将是一个函数。该函数应根据您的业务逻辑返回您已连接的直接下游任务的任务名称。这可能是下游的 1 到 N 个任务。下游任务 HAVE 没有要读取的内容,但是您可以使用 xcom 将元数据传递给它们。

def decide_which_path():
    if something is True:
        return "branch_a"
    else:
        return "branch_b"


branch_task = BranchPythonOperator(
    task_id='run_this_first',
    python_callable=decide_which_path,
    trigger_rule="all_done",
    dag=dag)

branch_task.set_downstream(branch_a)
branch_task.set_downstream(branch_b)

设置trigger_rule 很重要,否则将跳过所有其余部分,因为默认值为all_success

【讨论】:

  • trigger_rule 仍然如此吗?文档不建议您需要它,而只是一个虚拟任务,因为下游的其他任务(除了函数返回的任务)将被跳过airflow.incubator.apache.org/concepts.html#branching
  • 是的,那是正确的,所以这取决于下游任务是如何连接的。我想我假设所有的分支都合并回一个主线任务,但这可能甚至不是正常的用例(但这是我的正常用例)。
  • 您感兴趣的用例是什么?我目前的一个是'这个文件是否存在',如果它不存在,则继续创建它,否则虚拟任务然后 dag 成功退出。它专门用于将一些静态数据(永远不会改变)从 SQL 数据库加载到 hadoop。我希望它是幂等的,具有非常快速的无操作,如果不需要,则完全避免对源数据库的查询影响。
  • @punkrockpollly - 查看 XComs @airflow.apache.org/concepts.html(使用 Ctrl+F xcoms 可以快速找到它)XComs 非常适合小信号和数据大小。您还可以使用外部系统(如 redis 或数据库)来执行任务之间的通信。您的用例将有助于确定要走的路。
  • 非常感谢! It's important to set the trigger_rule or all the rest will be skipped, as the default is all_success 正是我所需要的!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-10-15
  • 2021-05-10
  • 2018-12-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多