YARN 架构与开发

LinuxBeginner
立即练习

简介

从 Hadoop 2.0 开始,Hadoop 引入了 YARN 这一全新的资源管理模式,提升了集群资源利用率,并实现了统一的资源管理和数据共享。本节基于已搭建好的 Hadoop 伪分布式集群,带你了解 YARN 框架的架构、工作原理、配置方法,以及开发和监控技巧。

本实验需要具备一定的 Java 编程基础。

开发步骤会提供完整的 Java 源文件。将每个完整文件粘贴到编辑器中,然后查看客户端如何提交应用,以及 ApplicationMaster 如何申请并启动任务容器。

YARN 架构与组件

本步骤将介绍 YARN 的架构及其组件的职责。

YARN 随 Hadoop 0.23 一同推出,是 MapReduce 2.0(MRv2)的一部分。它彻底改变了 Hadoop 集群中的资源管理和作业调度方式:

  • 拆分 JobTracker 的职责:MRv2 将 JobTracker 的功能拆分给不同的守护进程:ResourceManager 负责资源管理,ApplicationMaster 负责作业调度和监控。
  • 全局 ResourceManager:每个应用都对应一个 ApplicationMaster。应用可以是 MapReduce 作业,也可以是描述作业的 DAG。
  • 数据计算框架:ResourceManager、Slave 和 NodeManager 共同构成一个框架,由 ResourceManager 管理所有应用的资源。
  • ResourceManager 的组件:Scheduler 根据容量、队列等约束分配资源;ApplicationsManager 负责处理作业提交并运行 ApplicationMaster。
  • 资源分配:资源需求通过资源容器定义,其中包含内存、CPU、磁盘和网络等资源。
  • NodeManager 的职责:NodeManager 监控容器的资源使用情况,并向 ResourceManager 和 Scheduler 报告。
  • ApplicationMaster 的任务:ApplicationMaster 与 Scheduler 协商资源容器、跟踪状态并监控进度。

下图展示了这些组件之间的关系:

YARN 架构组件图

YARN 通过确保 API 与旧版本兼容,让正在运行的 MapReduce 任务可以平稳迁移。要在 Hadoop 集群中高效管理资源和调度作业,必须了解 YARN 的架构和组件。

启动 Hadoop 守护进程

本步骤将启动运行 YARN 应用所需的 Hadoop 守护进程。

在了解相关配置参数和 YARN 应用开发技巧之前,需要先启动 Hadoop 守护进程,以便随时使用 Hadoop。

先在桌面上双击打开 Xfce 终端,然后输入以下命令切换到 hadoop 用户:

su - hadoop

提示:hadoop 用户的密码是 hadoop。

切换完成后,即可启动 HDFS 和 YARN 框架相关的 Hadoop 守护进程。

在终端中输入以下命令启动这些守护进程:

/home/hadoop/hadoop/sbin/start-dfs.sh
/home/hadoop/hadoop/sbin/start-yarn.sh

启动完成后,可以运行 jps 命令检查相关守护进程是否正在运行。

hadoop:~$ jps
3378 NodeManager
3028 SecondaryNameNode
3717 Jps
2791 DataNode
2648 NameNode
3240 ResourceManager

准备配置文件

本步骤将介绍 Hadoop 的主要配置文件之一 yarn-site.xml,并查看如何通过该文件设置 YARN 集群。

为避免误改配置文件,最好先将 Hadoop 配置文件复制到其他目录,再打开副本查看。

在终端中输入以下命令,为配置文件创建一个新目录:

mkdir /home/hadoop/hadoop_conf

然后将 YARN 的主要配置文件 yarn-site.xml 从安装目录复制到新建的目录中。

在终端中输入以下命令完成复制:

cp /home/hadoop/hadoop/etc/hadoop/yarn-site.xml /home/hadoop/hadoop_conf/yarn-site.xml

然后使用 vim 编辑器打开该文件,查看其中的内容:

vim /home/hadoop/hadoop_conf/yarn-site.xml

配置文件的工作方式

本步骤将查看正在运行的 Hadoop 集群所使用的配置参数。

YARN 框架中有两个重要角色:ResourceManager 和 NodeManager。因此,文件中的配置项都用于设置这两个组件。

该文件可以设置许多配置项,但默认情况下不包含自定义配置项。例如,之前配置 Hadoop 伪分布式集群时指定了 aux-services 属性。当前打开的文件中只有这一项,如下所示:

hadoop:~$ cat /home/hadoop/hadoop/etc/hadoop/mapred-site.xml

...
<configuration>
    <property>
        <name>mapreduce.framework.name</name>
        <value>yarn</value>
    </property>
</configuration>

此配置项用于指定 NodeManager 上需要运行的依赖服务。这里指定的配置值是 mapreduce_shuffle,表示 MapReduce 程序在 YARN 上运行时需要使用默认值。

没有写入文件的配置项就不会生效吗?并非如此。如果没有在文件中明确指定配置参数,Hadoop 的 YARN 框架会读取内部文件中存储的默认值。在 yarn-site.xml 文件中明确指定的配置项会覆盖默认值,这样 Hadoop 系统就能适应不同的使用场景。

ResourceManager 配置项

要在 Hadoop 集群中高效管理资源和执行作业,必须了解并正确配置 yarn-site.xml 文件中的 ResourceManager 设置。以下是与 ResourceManager 相关的主要配置项:

  • **yarn.resourcemanager.address**:向客户端公开用于提交和终止应用的地址。默认端口为 8032。
  • **yarn.resourcemanager.scheduler.address**:向 ApplicationMaster 公开用于申请和释放资源的地址。默认端口为 8030。
  • **yarn.resourcemanager.resource-tracker.address**:向 NodeManager 公开用于发送心跳和获取任务的地址。默认端口为 8031。
  • **yarn.resourcemanager.admin.address**:向管理员公开用于执行管理命令的地址。默认端口为 8033。
  • **yarn.resourcemanager.webapp.address**:用于查看集群信息的 Web UI 地址。默认端口为 8088。
  • **yarn.resourcemanager.scheduler.class**:指定调度器的主类名称(例如 FIFO、CapacityScheduler、FairScheduler)。
  • 线程配置:
    • yarn.resourcemanager.resource-tracker.client.thread-count
    • yarn.resourcemanager.scheduler.client.thread-count
  • 资源分配:
    • yarn.scheduler.minimum-allocation-mb
    • yarn.scheduler.maximum-allocation-mb
    • yarn.scheduler.minimum-allocation-vcores
    • yarn.scheduler.maximum-allocation-vcores
  • NodeManager 管理:
    • yarn.resourcemanager.nodes.exclude-path
    • yarn.resourcemanager.nodes.include-path
  • 心跳配置:
    • yarn.resourcemanager.nodemanagers.heartbeat-interval-ms

配置这些参数可以调整 Hadoop 集群中 ResourceManager 的行为、资源分配、线程处理、NodeManager 管理方式和心跳间隔。了解这些配置项有助于避免问题,确保集群平稳运行。

NodeManager 配置项

要在 Hadoop 集群中高效管理资源和任务,必须正确配置 yarn-site.xml 文件中的 NodeManager 设置。以下是与 NodeManager 相关的主要配置项:

  • **yarn.nodemanager.resource.memory-mb**:指定 NodeManager 可用的物理内存总量。在 YARN 运行期间,此值保持不变。
  • **yarn.nodemanager.vmem-pmem-ratio**:设置虚拟内存与物理内存分配量的比率。默认值为 2.1。
  • **yarn.nodemanager.resource.cpu-vcores**:指定 NodeManager 可用的虚拟 CPU 总数。默认值为 8。
  • **yarn.nodemanager.local-dirs**:指定 NodeManager 存储中间结果的路径,可以配置多个目录。
  • **yarn.nodemanager.log-dirs**:指定 NodeManager 的日志目录路径,可以配置多个目录。
  • **yarn.nodemanager.log.retain-seconds**:指定 NodeManager 日志的最长保留时间。默认值为 10800 秒(3 小时)。

配置这些参数可以调整 Hadoop 集群中 NodeManager 的资源分配、内存管理、目录路径和日志保留设置,从而优化性能和资源利用率。了解这些配置项有助于确保集群平稳运行并高效执行任务。

查询配置项与默认配置参考

要查看 YARN 以及其他常见 Hadoop 组件的所有可用配置项,可以参考 Apache Hadoop 提供的默认配置文件。以下链接可用于查看默认配置:

查看这些默认配置文件,可以了解每个配置项的详细说明和用途,从而理解各个参数在 Hadoop 架构设计中的作用。

查看完配置后,关闭 vim 编辑器,结束对 Hadoop 配置设置的查看。

创建项目目录和文件

本步骤将创建应用的源文件。我们要构建一个完整但精简的 YARN 应用:客户端提交 ApplicationMaster,由 ApplicationMaster 申请一个任务容器并在其中打印问候语。

先创建项目目录。在终端中输入以下命令:

mkdir /home/hadoop/yarn_app

然后在项目目录中创建两个源代码文件。

第一个文件是 Client.java。在终端中使用 touch 命令创建该文件:

touch /home/hadoop/yarn_app/Client.java

然后创建 ApplicationMaster.java 文件:

touch /home/hadoop/yarn_app/ApplicationMaster.java
hadoop:~$ tree /home/hadoop/yarn_app/
/home/hadoop/yarn_app/
├── ApplicationMaster.java
└── Client.java

0 directories, 2 files

编写客户端代码

本步骤将编写用于提交应用的完整客户端。继续以 hadoop 用户身份操作。打开源文件:

vim /home/hadoop/yarn_app/Client.java

用下面的代码替换整个文件,包括 package 声明和 import 语句。在 Vim 中,先输入 :set paste 并按 Enter,然后按 i 进入插入模式。粘贴完整文件,等所有代码行显示出来后,按 Esc,再输入 :wq 保存并退出。

package com.labex.yarn.app;

import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.yarn.api.ApplicationConstants;
import org.apache.hadoop.yarn.api.records.*;
import org.apache.hadoop.yarn.client.api.YarnClient;
import org.apache.hadoop.yarn.client.api.YarnClientApplication;
import org.apache.hadoop.yarn.conf.YarnConfiguration;
import org.apache.hadoop.yarn.util.ConverterUtils;

public class Client {
    public static void main(String[] args) throws Exception {
        if (args.length != 1) {
            throw new IllegalArgumentException("Usage: Client /absolute/path/to/yarn-app.jar");
        }
        YarnConfiguration conf = new YarnConfiguration();
        YarnClient client = YarnClient.createYarnClient();
        client.init(conf);
        client.start();
        try {
            YarnClientApplication application = client.createApplication();
            ApplicationSubmissionContext context = application.getApplicationSubmissionContext();
            ApplicationId id = context.getApplicationId();
            FileSystem fs = FileSystem.get(conf);
            Path destination = new Path(fs.getHomeDirectory(), "yarn-app/" + id + "/app.jar");
            fs.mkdirs(destination.getParent());
            fs.copyFromLocalFile(new Path(args[0]), destination);
            FileStatus status = fs.getFileStatus(destination);
            LocalResource jar = LocalResource.newInstance(
                ConverterUtils.getYarnUrlFromPath(fs.makeQualified(destination)),
                LocalResourceType.FILE, LocalResourceVisibility.APPLICATION,
                status.getLen(), status.getModificationTime());
            Map<String, LocalResource> resources = new HashMap<>();
            resources.put("app.jar", jar);
            Map<String, String> environment = new HashMap<>();
            // The single-node lab uses the same Hadoop installation in every container.
            environment.put("CLASSPATH", "./app.jar:" + System.getProperty("java.class.path"));
            String command = ApplicationConstants.Environment.JAVA_HOME.$$()
                + "/bin/java -Xmx128m com.labex.yarn.app.ApplicationMaster"
                + " 1>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stdout"
                + " 2>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stderr";
            ContainerLaunchContext launch = ContainerLaunchContext.newInstance(
                resources, environment, Collections.singletonList(command), null, null, null);
            context.setApplicationName("LabEx YARN Hello");
            context.setQueue("default");
            context.setResource(Resource.newInstance(256, 1));
            context.setAMContainerSpec(launch);
            client.submitApplication(context);
            System.out.println("Application ID: " + id);
            long deadline = System.currentTimeMillis() + 180000;
            while (System.currentTimeMillis() < deadline) {
                ApplicationReport report = client.getApplicationReport(id);
                YarnApplicationState state = report.getYarnApplicationState();
                if (state == YarnApplicationState.FINISHED
                    || state == YarnApplicationState.FAILED
                    || state == YarnApplicationState.KILLED) {
                    System.out.println("Final status: " + report.getFinalApplicationStatus());
                    if (report.getFinalApplicationStatus() != FinalApplicationStatus.SUCCEEDED) {
                        throw new IllegalStateException(report.getDiagnostics());
                    }
                    return;
                }
                Thread.sleep(1000);
            }
            client.killApplication(id);
            throw new IllegalStateException("Application timed out after three minutes");
        } finally {
            client.stop();
        }
    }
}

YarnClient 会连接 ResourceManager 并获取应用 ID。客户端将 JAR 复制到 HDFS,并将其声明为本地资源 app.jar,这样 YARN 就能将它放入 ApplicationMaster 的工作目录。启动上下文提供 classpath 和 Java 命令;提交上下文则指定队列和容器资源。

classpath 使用了单节点实验中各容器共用的已安装 Hadoop 库。无需通过 Maven 或 Gradle 下载依赖。本示例适用于实验中未启用 Kerberos 的集群。

提交应用后,客户端会轮询应用报告,直到应用进入终止状态,并检查最终状态。只有应用运行时间超过三分钟时,才会调用 killApplication。最后一个步骤中,我们将编译并运行这两个类。

编写 ApplicationMaster 代码

本步骤将编写完整的 ApplicationMaster。它会向 ResourceManager 注册、申请一个容器,通过 NodeManager 启动问候任务,并报告任务结果。

以 hadoop 用户身份打开文件:

vim /home/hadoop/yarn_app/ApplicationMaster.java

用下面的代码替换整个文件。先输入 :set paste 并按 Enter,然后按 i。粘贴完整文件,等所有代码行显示出来后,按 Esc,再输入 :wq 保存并退出。

package com.labex.yarn.app;

import java.util.Collections;
import org.apache.hadoop.yarn.api.ApplicationConstants;
import org.apache.hadoop.yarn.api.protocolrecords.AllocateResponse;
import org.apache.hadoop.yarn.api.records.*;
import org.apache.hadoop.yarn.client.api.AMRMClient;
import org.apache.hadoop.yarn.client.api.NMClient;
import org.apache.hadoop.yarn.conf.YarnConfiguration;

public class ApplicationMaster {
    public static void main(String[] args) throws Exception {
        YarnConfiguration conf = new YarnConfiguration();
        AMRMClient<AMRMClient.ContainerRequest> rm = AMRMClient.createAMRMClient();
        NMClient nm = NMClient.createNMClient();
        rm.init(conf);
        nm.init(conf);
        rm.start();
        nm.start();
        try {
            rm.registerApplicationMaster("", 0, "");
            AMRMClient.ContainerRequest request = new AMRMClient.ContainerRequest(
                Resource.newInstance(256, 1), null, null, Priority.newInstance(0));
            rm.addContainerRequest(request);
            boolean launched = false;
            long deadline = System.currentTimeMillis() + 120000;
            while (System.currentTimeMillis() < deadline) {
                // Each allocate call also sends a heartbeat to the ResourceManager.
                AllocateResponse response = rm.allocate(launched ? 0.5f : 0.0f);
                for (Container container : response.getAllocatedContainers()) {
                    if (launched) {
                        rm.releaseAssignedContainer(container.getId());
                        continue;
                    }
                    rm.removeContainerRequest(request);
                    String command = "/bin/echo Hello-from-YARN"
                        + " 1>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stdout"
                        + " 2>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stderr";
                    ContainerLaunchContext launch = ContainerLaunchContext.newInstance(
                        Collections.emptyMap(), Collections.emptyMap(),
                        Collections.singletonList(command), null, null, null);
                    nm.startContainer(container, launch);
                    launched = true;
                }
                for (ContainerStatus status : response.getCompletedContainersStatuses()) {
                    if (status.getExitStatus() != 0) {
                        throw new IllegalStateException(status.getDiagnostics());
                    }
                    rm.unregisterApplicationMaster(FinalApplicationStatus.SUCCEEDED,
                        "Hello container completed", "");
                    return;
                }
                Thread.sleep(1000);
            }
            rm.unregisterApplicationMaster(FinalApplicationStatus.FAILED,
                "No successful container within two minutes", "");
            throw new IllegalStateException("Container timed out");
        } finally {
            nm.stop();
            rm.stop();
        }
    }
}

AMRMClient 负责注册和申请资源。每次调用 allocate 都会发送心跳,并返回新分配的容器和已完成容器的状态。NMClient 会在分配到的容器中启动命令。这个同步轮询循环让示例可以独立运行,无需依赖未定义的回调类或辅助方法。

任务会在容器日志中打印 Hello-from-YARN。如果任务退出状态为 0,ApplicationMaster 就会以 SUCCEEDED 状态注销。随后,客户端会报告该最终状态。ResourceManager 可能会将申请的 256 MB 向上调整到最小分配值;等待期间,循环会继续发送心跳。

应用启动流程

本步骤将编译并运行你编写的两个 Java 类,然后在 ResourceManager 的 Web 界面中查看应用。继续以 hadoop 用户身份在终端中操作。

编译并启动应用

切换到源文件目录:

cd /home/hadoop/yarn_app

创建输出目录:

mkdir -p classes

使用虚拟机中已安装的 Hadoop 库编译两个完整的源文件。--release 8 会生成与 Hadoop 的 Java 8 运行时兼容的类,即使默认编译器版本较新也不受影响:

javac --release 8 -cp "$(/home/hadoop/hadoop/bin/hadoop classpath --glob)" -d classes Client.java ApplicationMaster.java

如果看到关于弃用 API 的提示,可以忽略;编译过程中不应出现错误。将编译后的类打包:

jar cf yarn-app.jar -C classes .

提交客户端并保存其输出。最后一个参数是客户端要上传到 HDFS、供 ApplicationMaster 使用的 JAR:

set -o pipefail
/home/hadoop/hadoop/bin/yarn jar /home/hadoop/yarn_app/yarn-app.jar com.labex.yarn.app.Client /home/hadoop/yarn_app/yarn-app.jar | tee /home/hadoop/yarn_app/application.log

等待应用运行完成。除 Hadoop 日志消息外,预期还会看到:

Application ID: application_<timestamp>_<sequence>
Final status: SUCCEEDED

每次运行生成的 ID 都会不同。问候语会写入任务容器的日志,客户端则会打印应用 ID 和最终结果。本次运行使用的是你自己的客户端和 ApplicationMaster,而不是单独的预构建 MapReduce 示例。

查看应用运行结果

在桌面上打开 Firefox 并访问:

http://localhost:8088

在应用列表中找到 LabEx YARN Hello,并根据终端中打印的应用 ID 确认该应用。其状态应为 FINISHED,最终状态应为 SUCCEEDED。点击应用 ID 可查看详细信息。ResourceManager 负责跟踪应用;ApplicationMaster 管理任务容器;NodeManager 执行该容器。

总结

本实验基于已搭建好的 Hadoop 伪分布式集群,继续介绍 YARN 框架的架构、工作原理、配置方法,以及开发和监控技巧。课程提供了许多代码和配置文件,请仔细阅读。