Data Engineering Zoomcamp 实战:在 Hadoop YARN 上以伪分布式模式运行 Spark 集群
2026/9/21 15:46:05 网站建设 项目流程

Data Engineering Zoomcamp 实战:在 Hadoop YARN 上以伪分布式模式运行 Spark 集群

【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp

本篇指南是 Data Engineering Zoomcamp 第 6 模块(Batch Processing)环境搭建的关键一环,讲解如何为 Spark 作业安装 Hadoop 3.2.3、配置免密 SSH、以伪分布式(pseudo-distributed)模式启动 YARN,并让 Spark 以master=yarn提交作业。读完本文你将掌握:单机 YARN 的完整搭建流程、Spark 与 YARN 的对接配置、GCS 连接器接入,以及通过 Docker 容器运行时在 YARN 上提交 PySpark 作业(如本仓库的06_spark_sql.py)的完整方法。

为什么 Spark 模块需要 YARN

在 Data Engineering Zoomcamp 的 Spark 与 Docker 章节中,作业需要调度到一个集群资源管理器上执行,而这个角色由 YARN(Yet Another Resource Negotiator)承担。YARN 随 Hadoop 一起分发,因此即使你的 Spark 发行版本身自带(PySpark 4.x 通过pip/uv安装时会捆绑一份 Spark),你仍然需要单独安装 Hadoop 才能获得 YARN。

本仓库的安装文档(06-batch/setup/linux.md、06-batch/setup/macos.md、06-batch/setup/windows.md)解决的是 Spark 本身的安装与验证(local[*]模式);而本文档解决的是从单机本地模式升级到伪分布式集群模式的这一步。文档的默认假设是 Linux 环境,Windows 用户建议使用 WSL,macOS 上理论上也可运行。

与完全分布式(multi-node)不同,伪分布式模式将 NameNode、DataNode、ResourceManager、NodeManager 等全部进程运行在同一台机器上,既能真实体验 YARN 的调度与资源管理语义,又不需要多台服务器,适合课程学习和本地开发。

前置条件:免密 SSH 到 localhost

YARN 启动时会通过 SSH 连接到本机的各个节点,因此你必须能够无需输入密码地执行:

ssh localhost

如果该命令仍提示输入密码,需要把你的公钥追加到本机授权列表中:

cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys chmod 0600 ~/.ssh/authorized_keys

以上命令假设~/.ssh目录下已存在id_rsa.pub(若没有,先用ssh-keygen -t rsa生成)。chmod 0600用于收紧authorized_keys的权限,避免 SSH 因权限过宽而拒绝读取。

在 WSL 环境下,SSH 服务默认可能没有启动,需要手动启动:

sudo service ssh start

下载 Hadoop 3.2.3 二进制包

本模块使用的 Spark 预期对接 Hadoop 3.2 版本,因此仓库文档指定安装 Hadoop 3.2.3。可以从 Apache 官方镜像站获取离你最近的镜像并下载:

wget https://dlcdn.apache.org/hadoop/common/hadoop-3.2.3/hadoop-3.2.3.tar.gz

解压并进入目录:

tar xzfv hadoop-3.2.3.tar.gz cd hadoop-3.2.3/

解压后的目录结构即为 Hadoop 安装根目录(后续配置中的HADOOP_HOME就指向它),其下的etc/hadoop/存放核心配置文件,sbin/存放启动脚本。

在单节点上启动 YARN

YARN 的守护进程依赖 Java,第一步是在 Hadoop 的环境配置脚本中写入JAVA_HOME

echo "export JAVA_HOME=${JAVA_HOME}" >> etc/hadoop/hadoop-env.sh

说明:${JAVA_HOME}需要在你的 shell 中已正确设置(安装 Java 的步骤见 06-batch/setup/linux.md)。将JAVA_HOME写入hadoop-env.sh,可确保通过./sbin/脚本启动 YARN 时能找到 JVM。

随后启动 YARN:

./sbin/start-yarn.sh

启动成功后,YARN 的 ResourceManager Web UI 会监听在8088 端口,通过浏览器访问验证:

http://localhost:8088/

在该页面上可以看到集群的活跃节点(NodeManager)、已提交的应用列表以及资源使用情况——这是验证 YARN 是否正常工作的最直接手段。

让 Spark 以 YARN 作为 Master 运行

要在 YARN 上提交 Spark 作业,SparkSession 或spark-submit需要指定master="yarn"。同时,Spark 必须知道去哪里读取 YARN 的配置文件(如core-site.xmlyarn-site.xml),因此需要设置两个环境变量:

export HADOOP_HOME="${HOME}/spark/hadoop-3.2.3" export YARN_CONF_DIR="${HADOOP_HOME}/etc/hadoop"
  • HADOOP_HOME:指向你解压 Hadoop 的位置(${HOME}/spark/hadoop-3.2.3是仓库文档约定的路径,可按实际解压目录调整);
  • YARN_CONF_DIR:指向 Hadoop 的配置目录,Spark 从这里加载 Hadoop/YARN 的配置信息。

设置完成后,即可启动 Jupyter(在 notebook 中构建SparkSession.builder.master("yarn")),或直接使用spark-submit提交作业。

以本仓库的 06-batch/code/06_spark_sql.py 为例,该脚本通过argparse接收--input_green--input_yellow--output三个必填参数,读取纽约出租车 green/yellow 的 Parquet 数据,用 Spark SQL 按区域、月份、服务类型聚合收入指标后写出 Parquet。这种"读取 → 注册临时表 → Spark SQL 聚合 → 写出"的脚本结构,正是后续以spark-submit提交到 YARN 的标准形态。

连接 Spark 与 YARN 到 Google Cloud Storage(GCS)

当输入/输出数据存放在 GCS 时(如gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/...),需要让 Hadoop/YARN 认识gs://协议。做法是下载 GCS 连接器(GCS Connector for Hadoop 3):

gsutil cp gs://hadoop-lib/gcs/gcs-connector-hadoop3-2.2.5.jar .

随后需要修改两类配置文件:

  1. ${SPARK_HOME}/conf/spark-defaults.conf—— 仓库提供了参考模板 06-batch/setup/config/spark-defaults.conf:
spark-master yarn spark.hadoop.google.cloud.auth.service.account.enable true spark.hadoop.google.cloud.auth.service.account.json.keyfile /home/alexey

第一行直接指定默认 master 为yarn;后两行以spark.hadoop.为前缀,将属性透传给 Hadoop 配置,启用服务账号认证并指向凭证文件(示例路径为/home/alexey,请替换为你自己的凭证 JSON 文件路径)。

  1. ${YARN_CONF_DIR}/core-site.xml—— 仓库提供了完整参考配置 06-batch/setup/config/core-site.xml:
<?xml version="1.0" encoding="UTF-8"?> <?xml-stylesheet type="text/xsl" href="configuration.xsl"?> <configuration> <property> <name>fs.AbstractFileSystem.gs.impl</name> <value>com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS</value> </property> <property> <name>fs.gs.impl</name> <value>com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem</value> </property> <property> <name>fs.gs.auth.service.account.json.keyfile</name> <value>/home/alexey/.google/credentials/google_credentials.json</value> </property> <property> <name>fs.gs.auth.service.account.enable</name> <value>true</value> </property> </configuration>

其中前两条属性把gs://文件系统映射到 GCS 连接器的实现类;后两条启用服务账号认证并指定 Google 凭证 JSON 的本地路径(需替换为你自己的路径)。

core-site.xml中新增属性时,遵循如下的 property 模板即可:

<property> <name></name> <value></value> </property>

同样地,在 06-batch/code/cloud.md 中可以看到配套的上传与提交流程:先用gsutil -m cp -r pq/ gs://dtc_data_lake_de-zoomcamp-nytaxi/pq把本地 Parquet 数据上传到 GCS,再下载gcs-connector-hadoop3-2.2.5.jarlib/目录供 Spark 使用。

在 YARN 上使用 Docker 容器运行时提交作业

YARN 3.x 支持将 ApplicationMaster 与 Executor 运行在 Docker 容器中,这正好与课程的 Docker 章节衔接。仓库文档给出的思路是:从 Hadoop 官方文档(hadoop-yarn-site/DockerContainers.html,对应 Hadoop 3.2.3)复制 Docker 容器相关的 YARN 配置,然后使用自定义镜像执行spark-submit

提交命令的核心部分如下:

MOUNTS="$HADOOP_HOME:$HADOOP_HOME:ro,/etc/passwd:/etc/passwd:ro,/etc/group:/etc/group:ro" IMAGE_ID="pyspark-docker:test" spark-submit \ --master yarn \ --conf spark.yarn.appMasterEnv.YARN_CONTAINER_RUNTIME_TYPE=docker \ --conf spark.yarn.appMasterEnv.YARN_CONTAINER_RUNTIME_DOCKER_IMAGE=${IMAGE_ID} \ --conf spark.executorEnv.YARN_CONTAINER_RUNTIME_TYPE=docker \ --conf spark.executorEnv.YARN_CONTAINER_RUNTIME_DOCKER_IMAGE=${IMAGE_ID} \ 06_spark_sql.py \ --input_green=gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/green/2021/*/ \ --input_yellow=gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/yellow/2021/*/ \ --output=gs://dtc_data_lake_de-zoomcamp-nytaxi/report-2021

要点拆解:

  • MOUNTS:以hostPath:containerPath:mode的格式把 Hadoop 安装目录、/etc/passwd/etc/group以只读(ro)方式挂载进容器,保证容器内能读取 Hadoop 配置且用户/组信息一致;
  • IMAGE_ID="pyspark-docker:test":指定自定义 Spark 镜像(仓库在 06-batch/setup/config/spark.dockerfile 中提供了最简 Dockerfile,仅一行FROM library/openjdk:11,可作为构建基础镜像的起点);
  • 通过--conf分别对ApplicationMasterspark.yarn.appMasterEnv.*)和Executorspark.executorEnv.*)声明容器运行时类型为 docker 及对应镜像;
  • 作业脚本06_spark_sql.py的三个参数分别指向 GCS 上的 green、yellow 数据以及聚合结果输出位置。

注意:gs://路径与 GCS 连接器配置(上一节)必须同时就位,Docker 容器运行时才能访问远端数据。

常见排查要点

  • ssh localhost仍需密码:检查authorized_keys权限是否为0600~/.ssh目录权限是否合理,WSL 用户记得先sudo service ssh start
  • start-yarn.sh报 Java 相关错误:确认JAVA_HOME已写入etc/hadoop/hadoop-env.sh,且路径指向有效的 JDK;
  • 8088 端口无法访问:确认 ResourceManager 进程已启动,且无防火墙拦截;
  • gs://路径读写失败:确认core-site.xml已配置fs.gs.impl与服务账号凭证,且spark-defaults.conf中的spark.hadoop.*前缀属性正确;
  • 提交作业时找不到 YARN 配置:确认YARN_CONF_DIR指向${HADOOP_HOME}/etc/hadoop

参考来源

本文档的核心依据为仓库中的 06-batch/setup/hadoop-yarn.md,配套的配置文件见 06-batch/setup/config/core-site.xml、06-batch/setup/config/spark-defaults.conf 与 06-batch/setup/config/spark.dockerfile。原始文档引用的外部资料包括 Hadoop 官方单机集群配置文档(hadoop-project-dist/hadoop-common/SingleCluster.html)与 Spark 自定义 Hadoop/Hive 配置文档(configuration.html#custom-hadoophive-configuration),前者对应本指南中的伪分布式 YARN 搭建,后者对应YARN_CONF_DIRspark.hadoop.*的配置原理,如需深入可自行查阅对应版本文档。

【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询