相关疑难解决方法(0)

在Airflow中创建动态工作流的正确方法

问题

在Airflow中是否有任何方法可以创建工作流程,以便任务数量B.*在任务A完成之前是未知的?我查看了子标记,但看起来它只能用于必须在Dag创建时确定的一组静态任务.

dag会触发工作吗?如果是这样,请你举个例子.

我有一个问题是,在任务A完成之前,无法知道计算任务C所需的任务B的数量.每个任务B.*将需要几个小时来计算,不能合并.

              |---> Task B.1 --|
              |---> Task B.2 --|
 Task A ------|---> Task B.3 --|-----> Task C
              |       ....     |
              |---> Task B.N --|
Run Code Online (Sandbox Code Playgroud)

想法#1

我不喜欢这个解决方案,因为我必须创建一个阻塞的ExternalTask​​Sensor,所有的任务B.*需要2到24小时才能完成.所以我认为这不是一个可行的解决方案.当然有一种更简单的方法吗?或者Airflow不是为此而设计的?

Dag 1
Task A -> TriggerDagRunOperator(Dag 2) -> ExternalTaskSensor(Dag 2, Task Dummy B) -> Task C

Dag 2 (Dynamically created DAG though python_callable in TriggerDagrunOperator)
               |-- Task B.1 --|
               |-- Task B.2 --|
Task Dummy A --|-- Task B.3 --|-----> Task Dummy B
               |     ....     |
               |-- Task B.N --|
Run Code Online (Sandbox Code Playgroud)

编辑1: …

python workflow airflow

66
推荐指数
8
解决办法
2万
查看次数

Airflow 任务能否在运行时动态生成 DAG?

我有一个不规则上传的上传文件夹。对于每个上传的文件,我想生成一个特定于该文件的 DAG。

我的第一个想法是使用 FileSensor 来执行此操作,该文件传感器监视上传文件夹,并以新文件的存在为条件,触发创建单独 DAG 的任务。从概念上讲:

Sensor_DAG (FileSensor -> CreateDAGTask)

|-> File1_DAG (Task1 -> Task2 -> ...)
|-> File2_DAG (Task1 -> Task2 -> ...)
Run Code Online (Sandbox Code Playgroud)

在我最初的实现中,CreateDAGTaskPythonOperator通过将它们放置在全局命名空间中来创建 DAG 全局变量(请参阅此 SO 答案),如下所示:

Sensor_DAG (FileSensor -> CreateDAGTask)

|-> File1_DAG (Task1 -> Task2 -> ...)
|-> File2_DAG (Task1 -> Task2 -> ...)
Run Code Online (Sandbox Code Playgroud)

主 DAG 然后通过一个调用这个逻辑PythonOperator

# File-sensing DAG
default_args = {
    "depends_on_past" : False,
    "start_date"      : datetime(2020, 7, 16),
    "retries"         : 1,
    "retry_delay"     : timedelta(hours=5),
} …
Run Code Online (Sandbox Code Playgroud)

airflow airflow-scheduler

7
推荐指数
1
解决办法
6257
查看次数

如何在 Airflow 中动态创建子标签

我有一个主 dag,它检索文件并将该文件中的数据拆分为单独的 csv 文件。我必须为这些 csv 文件的每个文件完成另一组任务。例如(上传到 GCS,插入到 BigQuery)如何根据文件数量为每个文件动态生成一个 SubDag?SubDag 将定义上传到 GCS、插入到 BigQuery、删除 csv 文件等任务)

所以现在,这就是它的样子

main_dag = DAG(....)
download_operator = SFTPOperator(dag = main_dag, ...)  # downloads file
transform_operator = PythonOperator(dag = main_dag, ...) # Splits data and writes csv files

def subdag_factory(): # Will return a subdag with tasks for uploading to GCS, inserting to BigQuery.
    ...
    ...
Run Code Online (Sandbox Code Playgroud)

如何为在 transform_operator 中生成的每个文件调用 subdag_factory?

airflow

5
推荐指数
1
解决办法
5379
查看次数

标签 统计

airflow ×3

airflow-scheduler ×1

python ×1

workflow ×1