Airflow DAG出现create_job_flow到remove_cluster意外依赖问题咨询
问题1解答
不会触发集群提前删除。Airflow中任务触发的前提是所有上游任务全部执行成功,当前remove_cluster有两个上游:你显式声明的alter_partitions、以及引用create_job_flow.output生成的隐式依赖create_job_flow。哪怕create_job_flow早就执行完成,只要alter_partitions还没运行成功,remove_cluster就不会被触发。
问题2解答
可以移除这条隐式依赖边,方法是把remove_cluster中job_flow_id的取值方式从直接引用create_job_flow.output改为手动用模板拉取XCom值,修改后的算子代码如下:
remove_cluster = EmrTerminateJobFlowOperator( task_id='remove_cluster', job_flow_id="{{ ti.xcom_pull(task_ids='create_job_flow', key='return_value') }}", aws_conn_id='aws_default', )
直接引用算子output属性时,Airflow会自动在两个算子间添加依赖;改用模板字符串手动指定XCom拉取逻辑时,不会自动生成隐式依赖,就能去掉这条多余的边。
问题3解答
是的,remove_cluster仍然会按照你设置的规则等待alter_partitions执行完成后再运行。正如问题1的解释,任务需要等所有上游完成才会触发,只要你保留了alter_partitions >> remove_cluster的显式依赖,哪怕存在额外的隐式依赖,也不会改变remove_cluster需要等待alter_partitions完成的执行逻辑。
内容的提问来源于stack exchange,提问作者seou1
相关产品推荐
相关产品推荐

