从HDFS到Spark On Yarn:在WSL里用VSCode完整跑通你的第一个大数据流水线

大数据技术栈的复杂性常常让初学者望而生畏——HDFS、Yarn、Spark这些名词听起来高大上,但如何让它们在实际开发中协同工作?本文将带你用WSL和VSCode搭建一个伪分布式环境,完成从数据存储到分布式计算的完整闭环。不同于碎片化的环境搭建教程,我们聚焦真实工作流:你会亲手实现HDFS文件操作、MapReduce经典案例、Spark应用提交到Yarn的全过程,最终获得"真正能用起来"的实战能力。

1. 环境准备:WSL2与基础服务配置

1.1 WSL2初始化配置

在Windows Terminal中执行以下命令完成基础环境部署:

wsl --install -d Ubuntu-22.04

安装完成后需要配置三个关键目录结构:

/opt
├── module    # 存放解压后的软件
├── software  # 存放安装包
└── data      # 实验数据集存放位置

提示:建议通过sudo chown -R $USER /opt修改目录所有权,避免后续操作频繁使用sudo

1.2 Hadoop单节点集群部署

下载Hadoop 3.3.6和JDK17后,环境变量配置应包含以下核心参数:

# ~/.bashrc追加内容
export JAVA_HOME=/opt/module/jdk-17.0.8
export HADOOP_HOME=/opt/module/hadoop-3.3.6
export PATH=$PATH:$JAVA_HOME/bin:$HADOOP_HOME/bin:$HADOOP_HOME/sbin

关键配置文件修改对比:

文件 必须配置项 伪分布式典型值
core-site.xml fs.defaultFS hdfs://localhost:9000
hdfs-site.xml dfs.replication 1
yarn-site.xml yarn.nodemanager.aux-services mapreduce_shuffle

启动服务时建议使用组合命令脚本:

#!/bin/bash
hdfs namenode -format && \
start-dfs.sh && \
start-yarn.sh

2. 数据流水线实践:从HDFS到MapReduce

2.1 HDFS文件系统操作实战

通过命令行与Web UI双视角操作HDFS:

# 创建用户目录
hadoop fs -mkdir -p /user/$(whoami)

# 上传本地文件到HDFS
echo "Hello World" > test.txt
hadoop fs -put test.txt input/

# 查看文件块信息
hadoop fsck /user/$(whoami)/input/test.txt -files -blocks

Web界面访问http://localhost:9870可直观查看文件分布。当遇到安全模式限制时,使用:

hdfs dfsadmin -safemode leave

2.2 经典WordCount实现剖析

MapReduce作业提交包含三个关键阶段:

  1. Mapper阶段:文本分割为<单词,1>键值对
  2. Shuffle阶段:相同单词聚合到同一Reducer
  3. Reducer阶段:统计每个单词出现次数

执行内置示例的完整命令:

hadoop jar $HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.3.6.jar \
wordcount input output

通过hadoop fs -cat output/part-r-00000查看结果时,注意输出目录必须不存在,否则会报错。

3. Spark On Yarn集成开发

3.1 Spark伪分布式配置

Spark与Yarn集成需要两个关键配置:

# spark-env.sh
export HADOOP_CONF_DIR=$HADOOP_HOME/etc/hadoop
export YARN_CONF_DIR=$HADOOP_HOME/etc/hadoop

# 测试资源调度是否正常
spark-submit --master yarn --deploy-mode client \
--class org.apache.spark.examples.SparkPi \
$SPARK_HOME/examples/jars/spark-examples_2.12-3.4.1.jar 10

常见问题排查表:

现象 可能原因 解决方案
ApplicationMaster启动失败 内存配置不足 调整yarn.scheduler.minimum-allocation-mb
Spark作业卡在ACCEPTED状态 资源竞争 检查yarn.resourcemanager.scheduler.class配置
日志中出现ClassNotFound 依赖包缺失 使用--jars参数指定额外依赖

3.2 PySpark开发环境搭建

使用Miniconda创建隔离环境:

conda create -n pyspark python=3.10
conda activate pyspark
pip install pyspark==3.4.1 pandas pyarrow

VSCode远程开发配置要点:

  1. 安装"Remote - WSL"扩展
  2. 连接WSL实例后选择Python解释器路径: /opt/module/miniconda3/envs/pyspark/bin/python
  3. 配置.vscode/settings.json
{
  "python.pythonPath": "/opt/module/miniconda3/envs/pyspark/bin/python",
  "python.linting.enabled": true
}

4. 完整数据流水线实战

4.1 从HDFS到Spark的ETL流程

实现一个包含数据清洗的完整PySpark作业:

# etl_pipeline.py
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("HDFS_ETL") \
    .getOrCreate()

# 从HDFS读取原始数据
raw_df = spark.read.text("hdfs://localhost:9000/user/input/logs.txt")

# 数据转换操作
cleaned_df = raw_df.filter("value IS NOT NULL") \
    .withColumn("timestamp", split(col("value"), " ")[0]) \
    .withColumn("message", substring_index(col("value"), " ", -3))

# 写回HDFS
cleaned_df.write.mode("overwrite") \
    .parquet("hdfs://localhost:9000/user/output/cleaned_logs")

提交作业时指定Executor资源:

spark-submit --master yarn \
--num-executors 2 \
--executor-cores 1 \
--executor-memory 1G \
etl_pipeline.py

4.2 性能优化技巧

通过Web UI(http://localhost:8088)观察作业运行情况时,重点关注:

  1. 数据倾斜:某些Task执行时间明显长于其他
  2. GC时间:垃圾回收耗时占比过高
  3. Shuffle数据量:跨节点数据传输规模

优化配置示例:

conf = SparkConf() \
    .set("spark.sql.shuffle.partitions", "200") \
    .set("spark.executor.extraJavaOptions", "-XX:+UseG1GC") \
    .set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")

5. 开发调试与问题排查

5.1 VSCode调试配置

.vscode/launch.json中添加Spark调试配置:

{
  "name": "PySpark Debug",
  "type": "python",
  "request": "launch",
  "program": "${file}",
  "env": {
    "PYSPARK_PYTHON": "/opt/module/miniconda3/envs/pyspark/bin/python",
    "PYSPARK_DRIVER_PYTHON": "/opt/module/miniconda3/envs/pyspark/bin/python"
  },
  "args": []
}

调试时可通过spark.sparkContext.uiWebUrl查看实时运行状态。

5.2 Yarn日志分析技巧

获取完整应用日志的命令:

yarn logs -applicationId application_123456789_0001

关键日志文件位置:

  • Driver日志yarn.nodemanager.log-dirs指定目录下的stdout/stderr
  • Executor日志:各节点container_*/目录下的日志文件

对于长时间运行的Spark Streaming应用,建议添加日志滚动配置:

logger = sc._jvm.org.apache.log4j
logger.LogManager.getLogger("org").setLevel(logger.Level.WARN)
logger.LogManager.getLogger("akka").setLevel(logger.Level.WARN)

更多推荐