【问题标题】:Spark Yarn Architecture火花纱架构
【发布时间】:2016-07-12 22:58:55
【问题描述】:

在我关注的教程中,我对这张图片有疑问。因此,在基于纱线的架构中,基于此图像,火花应用程序的执行看起来像这样:

首先,您有一个在客户端节点或某个数据节点上运行的驱动程序。在这个驱动程序中(类似于 java 中的驱动程序?)由您提交给 Spark 上下文的代码(用 java、python、scala 等编写)组成。然后,该 spark 上下文表示与 HDFS 的连接,并将您的请求提交给 Hadoop 生态系统中的资源管理器。然后资源管理器与 Name 节点通信,以确定集群中的哪些数据节点包含客户端节点请求的信息。火花上下文还将在将运行任务的工作节点上放置一个执行程序。然后节点管理器将启动执行器,该执行器将运行 Spark 上下文给它的任务,并将客户端从 HDFS 请求的数据返回给驱动程序。

以上解释正确吗?

由于 HDFS 中的数据在不同的数据节点上被复制了 3 次,驱动程序还会向每个数据节点发送三个执行程序以从 HDFS 检索数据吗?

【问题讨论】:

    标签: scala hadoop apache-spark hdfs


    【解决方案1】:

    您的解释接近现实,但您似乎在某些方面有些困惑。

    让我们看看我能否让你更清楚。

    假设您有 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

    【讨论】:

    • 非常感谢您的详细解释!关于资源管理器和名称节点如何协同工作以找到工作节点。所以基本上你的文件的三个副本存储在 HDFS 的三个不同的数据节点上。资源管理器将根据数据局部性选择具有第一个 HDFS 块的工作节点,并联系该工作节点上的 NodeManager 以创建一个纱线容器(JVM),在其中运行火花执行器。如果其他区块在这个“范围”内不可用,那么它会去其他工作节点,将其他区块转移过来
    • 到资源管理器最初找到的最近数据节点的网络(运行该 spark 执行器)正确吗?
    • 另外关于您在上面编写的示例字数统计程序中的输入文件是否来自 HDFS?
    • 为了解释我的示例,我假设它来自 hdfs,但相同的源代码适用于本地文件和 hdfs 文件。如果您使用 spark-submit,spark 将假定输入文件路径是相对于 hdfs 的,如果您在 Intellij idea 中将其作为 Java 程序运行,它将假定它是一个本地文件。如果从 Intellij 运行,您甚至可以使用 hdfs 文件,但在这种情况下,您必须指定 hdfs://。在最后一种情况下,由于您在笔记本电脑上运行并从远程 hdfs 集群读取数据,因此您将失去局部性
    • 哦,这很有道理,太棒了!感谢您的所有澄清,绝对有很大帮助!
    猜你喜欢
    • 1970-01-01
    • 2019-11-13
    • 1970-01-01
    • 2017-12-25
    • 2014-12-17
    • 2017-03-21
    • 2017-08-27
    • 2016-10-11
    • 2019-06-26
    相关资源
    最近更新 更多