gpt4 book ai didi

python - 想要创建当前任务下游的 Airflow 任务

转载 作者:太空宇宙 更新时间:2023-11-04 09:49:55 25 4
gpt4 key购买 nike

我对 Airflow 几乎是全新的。

我有两个步骤:

  1. 获取所有符合条件的文件
  2. 解压文件

文件压缩后为半个 gig,未压缩时为 2 - 3 gig。我一次可以轻松处理 20 多个文件,这意味着解压缩所有这些文件的运行时间比任何合理的超时都长

我可以使用 XCom 获得步骤 1 的结果,但我想做的是这样的:

def processFiles (reqDir, gvcfDir, matchSuffix):
theFiles = getFiles (reqDir, gvcfDir, matchSuffix)

for filePair in theFiles:
task = PythonOperator (task_id = "Uncompress_" + os.path.basename (theFile),
python_callable = expandFile,
op_kwargs = {'theFile': theFile},
dag = dag)
task.set_upstream (runThis)

问题是“runThis”是调用processFiles的PythonOperator,所以必须在processFiles之后声明。

有什么方法可以让它工作吗?

这是 XCom 存在的原因,我应该放弃这种方法并使用 XCom 吗?

最佳答案

关于您提出的解决方案,我认为您不能使用 XComs 来实现这一点,因为它们仅适用于实例,而不是在您定义 DAG 时可用(据我所知)。

但是您可以使用 SubDAG实现你的目标。 SubDagOperator 获取一个函数,该函数将在执行操作符时调用并生成 DAG,让您有机会动态创建工作流的子部分。

您可以使用这个简单的示例来测试这个想法,它会在每次调用时生成随机的任务:

import airflow
from builtins import range
from random import randint
from airflow.operators.bash_operator import BashOperator
from airflow.operators.subdag_operator import SubDagOperator
from airflow.models import DAG

args = {
'owner': 'airflow',
'start_date': airflow.utils.dates.days_ago(2)
}

dag = DAG(dag_id='dynamic_dag', default_args=args)

def generate_subdag(parent_dag, dag_id, default_args):
# pseudo-randomly determine a number of tasks to be created
n_tasks = randint(1, 10)

subdag = DAG(
'%s.%s' % (parent_dag.dag_id, dag_id),
schedule_interval=parent_dag.schedule_interval,
start_date=parent_dag.start_date,
default_args=default_args
)
for i in range(n_tasks):
i = str(i)
task = BashOperator(task_id='echo_%s' % i, bash_command='echo %s' % i, dag=subdag)

return subdag

subdag_dag_id = 'dynamic_subdag'

SubDagOperator(
subdag=generate_subdag(dag, subdag_dag_id, args),
task_id=subdag_dag_id,
dag=dag
)

如果您执行此操作,您会注意到在不同的运行中,SubDAG 可能包含不同数量的任务(我使用 1.8.0 版对此进行了测试)。您可以通过访问图 TableView 、单击灰色的 SubDAG 节点然后单击“放大 SubDAG”来访问 WebUI 上的 SubDAG View 。

您可以通过列出文件并为每个文件创建一个任务来使用此概念,而不是像示例中那样仅以随机数生成它们。任务本身可以并行排列(就像我所做的那样)、顺序排列或以任何有效的定向非循环布局排列。

关于python - 想要创建当前任务下游的 Airflow 任务,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/48197709/

25 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com