【问题标题】:How to control concurrency of a subdag or taskgroup in Airflow?如何控制 Airflow 中 subdag 或 taskgroup 的并发性?
【发布时间】:2021-07-08 13:02:39
【问题描述】:

我有一个简单的 ETL 工作流程,由我要提取的每个表的三个任务组成:

1_extract_to_tmp >> 2_push_to_s3 >> 3_delete_tmp

由于我想为多个表重现相同的步骤,我正在考虑将这些任务分组到一个 TaskGroup 或一个 subDAG 中,并为每个要提取的表动态实例化它。我的最终 DAG 将如下所示:

根据我的阅读,TaskGroup 是现在首选的解决方案。

问题是我还需要控制并行提取的数量,因为我没有足够的本地磁盘空间来存储所有数据,并且显然 subDag 运算符不符合池设置并且 TaskGroup 没有.

您知道实现此目的的方法吗?

我做错了吗?

也许我的 DAG 不应该这样设计。

【问题讨论】:

    标签: concurrency airflow pool


    【解决方案1】:

    TaskGroup 只是一个 UI 功能。它并不真正包含逻辑,因此无论使用 TaskGroup,您都可以执行任何并行限制。

    你没有解释你到底想限制什么。 如果您想限制可以在您的 dag 上并行运行的总体任务(覆盖 airflow.cfg 默认值),请在您的 DAG 承包商中设置 concurrency

    dag = DAG(dag_id='my_dag', concurrency=5, ...)
    

    如果您希望限制并发使用以保护特定资源(数据库、api 等),请使用pools

    my_op = PythonOperator(python_callable=func,
                              task_id='my_task',
                              pool='my_pool',
                              dag=dag)
    

    只有在池中有空闲槽时才会执行任务。可用槽的数量是在您创建池时定义的。

    至于 SubDags,我根本不推荐使用它。虽然它没有被正式弃用,但它可能会被弃用。见https://github.com/apache/airflow/issues/12292

    【讨论】:

    • 我在原始问题中添加了 DAG 图。我希望它能澄清一些事情。在图中,Extract(Table1)、Extract(Table2)、Extract(Table3) 和 Extract(Table4) 都是 TaskGroup。它们可以并行运行,但我想限制并发级别,以便只允许 2 个任务组同时运行。看起来不太可能。
    • 答案是一样的。使用池:)
    猜你喜欢
    • 2019-10-15
    • 1970-01-01
    • 2022-12-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-07-17
    相关资源
    最近更新 更多