您的解释接近现实,但您似乎在某些方面有些困惑。
让我们看看我能否让你更清楚。
假设您有 Scala 中的字数统计示例。
object WordCount {
def main(args: Array[String]) {
val inputFile = args(0)
val outputFile = args(1)
val conf = new SparkConf().setAppName("wordCount")
val sc = new SparkContext(conf)
val input = sc.textFile(inputFile)
val words = input.flatMap(line => line.split(" "))
val counts = words.map(word => (word, 1)).reduceByKey{case (x, y) => x + y}
counts.saveAsTextFile(outputFile)
}
}
在每个 spark 作业中,您都有一个初始化步骤,在该步骤中您创建一个 SparkContext 对象,提供一些配置,例如 appname 和 master,然后您读取一个 inputFile,处理它并将处理结果保存在磁盘上。所有这些代码都在驱动程序中运行,除了进行实际处理的匿名函数(传递给 .flatMap、.map 和 reduceByKey 的函数)以及在集群上远程运行的 I/O 函数 textFile 和 saveAsTextFile。
这里的 DRIVER 是在您使用spark-submit 提交代码的同一节点上本地运行的程序部分的名称(在您的图片中称为客户端节点)。只要您对 YARN 集群具有 spark-submit 和网络访问权限,就可以从任何机器(ClientNode、WorderNode 甚至 MasterNode)提交您的代码。为简单起见,我假设 Client 节点是您的笔记本电脑,Yarn 集群由远程机器组成。
为简单起见,我将省略 Zookeeper,因为它用于为 HDFS 提供高可用性,并且不涉及运行 spark 应用程序。不得不提一下,Yarn Resource Manager 和 HDFS Namenode 是 Yarn 和 HDFS 中的角色(实际上它们是在 JVM 中运行的进程),它们可以存在于同一个主节点上,也可以存在于不同的机器上。即使是 Yarn 节点管理器和数据节点也只是角色,但它们通常位于同一台机器上以提供数据局部性(处理靠近数据存储的位置)。
当您提交您的应用程序时,您首先联系资源管理器,它与 NameNode 一起尝试查找可用于运行 Spark 任务的 Worker 节点。为了利用数据局部性原则,资源管理器将优先选择在同一台机器上存储您必须处理的文件的 HDFS 块(每个块的 3 个副本中的任何一个)的工作节点。如果没有具有这些块的工作节点可用,它将使用任何其他工作节点。在这种情况下,由于本地数据不可用,HDFS 块必须通过网络从任何数据节点移动到运行 spark 任务的节点管理器。这个过程是针对每个生成文件的块完成的,所以有些块可以在本地找到,有些则必须移动。
当 ResourceManager 找到一个可用的工作节点时,它将联系该节点上的 NodeManager 并要求它创建一个 Yarn Container (JVM) 来运行 spark executor。在其他集群模式(Mesos 或 Standalone)中,您不会有 Yarn 容器,但 spark executor 的概念是相同的。 spark executor 作为 JVM 运行,可以运行多个任务。
在客户端节点上运行的驱动程序和在 spark 执行器上运行的任务保持通信以运行您的作业。如果驱动程序正在您的笔记本电脑上运行并且您的笔记本电脑崩溃,您将失去与任务的连接并且您的工作将失败。这就是为什么当 spark 在 Yarn 集群中运行时,您可以指定是在笔记本电脑上运行驱动程序“--deploy-mode=client”还是在纱线集群上作为另一个纱线容器“--deploy-mode=cluster ”。更多详情请查看spark-submit