从HDFS到Spark On Yarn:在WSL里用VSCode完整跑通你的第一个大数据流水线
从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作业提交包含三个关键阶段:
- Mapper阶段:文本分割为<单词,1>键值对
- Shuffle阶段:相同单词聚合到同一Reducer
- 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远程开发配置要点:
- 安装"Remote - WSL"扩展
- 连接WSL实例后选择Python解释器路径:
/opt/module/miniconda3/envs/pyspark/bin/python - 配置
.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)观察作业运行情况时,重点关注:
- 数据倾斜:某些Task执行时间明显长于其他
- GC时间:垃圾回收耗时占比过高
- 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)
更多推荐



所有评论(0)