Yan*_*mer 7 airflow databricks
我正在使用气流来触发数据块上的作业。我有许多运行数据块作业的 DAG,我希望只使用一个集群而不是多个集群,因为据我所知,这将降低这些任务产生的成本。
使用DatabricksSubmitRunOperator有两种方法可以在数据块上运行作业。使用正在运行的集群通过 id 调用它
'existing_cluster_id' : '1234-567890-word123',
Run Code Online (Sandbox Code Playgroud)
或者启动一个新集群
'new_cluster': {
'spark_version': '2.1.0-db3-scala2.11',
'num_workers': 2
},
Run Code Online (Sandbox Code Playgroud)
现在我想尽量避免为每个任务启动一个新集群,但是集群在停机期间关闭,因此它不再通过它的 id 可用,我会得到一个错误,所以我认为唯一的选择是新集群。
1) 有没有办法让集群即使在关闭时也可以通过 id 调用?
2)人们是否只是让集群保持活力?
3)还是我完全错了,为每个任务启动集群不会产生更多成本?
4)有什么我完全错过的吗?
基于@YannickSSE\'s 评论响应的更新
\n我不使用databricks;您是否可以使用与您可能或可能不希望正在运行的集群相同的 ID 启动一个新集群,并使其在运行时处于无操作状态?也许不会,或者你可能不会问这个。响应:否,启动新集群时无法提供 id。
你能编写一个 python 或 bash 运算符来测试集群是否存在吗?(响应:这将是一个测试作业提交\xe2\x80\xa6,不是最好的方法。)如果找到它并成功,下游任务将使用现有集群 ID 触发您的作业,但如果没有找到另一个下游任务任务可以使用来执行相同的任务,但使用新的集群。那么这两个任务都可以有一个下游任务。(响应:或者使用分支运算符来确定执行的运算符。)trigger_rule all_failedDatabricksSubmitRunOperatortrigger_rule one_success
这可能并不理想,因为我想您的集群 ID 会不时发生变化,导致您必须跟上。\xe2\x80\xa6 集群是该操作员的 databricks hook 连接的一部分,并且可以更新吗?也许您想在需要它的任务中指定它{{ var.value.<identifying>_cluster_id }},并将其作为气流变量进行更新。(响应:集群 ID 不在钩子中,因此变量或 DAG 文件每当发生变化时都必须更新。)