Spark 官方示例学习:测试数据生成与实战代码
基于 Apache Spark 官方示例体系,结合阿里云 EMR 集群环境,生成对应测试数据并提供可直接运行的 PySpark 实现,覆盖官方入门四大经典场景:WordCount 词频统计、Pi 圆周率计算、Spark SQL 结构化查询、PageRank 网页排名。所有实现对齐官方逻辑,仅做云环境适配与中文注解。
前置说明
- 运行环境:阿里云 EMR-5.21.0 集群,Spark 3.5.3,YARN 资源调度,HDFS 分布式存储
- 数据路径:所有测试数据默认存放在 HDFS 的
/test/目录下 - 运行方式:同时提供交互式(pyspark 终端逐行调试)和脚本提交(spark-submit 正式运行)两种方式
- 学习路径:从基础文本计算 → 数值计算 → 结构化查询 → 图算法,循序渐进对齐官方示例体系
一、文本类测试数据 + WordCount 词频统计
WordCount 是 Spark 官方入门的经典示例,对应官方 Quick Start、RDD 编程指南中的核心案例,也是 MapReduce 编程范式的代表。
1.1 生成测试数据
SSH 登录 Master 节点,执行以下命令生成模拟文本数据并上传至 HDFS:
# 生成测试文本:模拟多段英文内容,覆盖不同高频单词
cat > words.txt << EOF
hello spark hello hadoop hello emr
spark is fast hadoop is stable
emr is easy to use spark powerful
hadoop spark emr spark hadoop
spark sql spark streaming spark mllib
hello world hello spark hello hadoop
big data spark hadoop emr
spark on yarn spark on k8s
hdfs mapreduce yarn spark
data analysis machine learning spark
EOF
# 创建 HDFS 测试目录并上传文件
hdfs dfs -mkdir -p /test
hdfs dfs -put words.txt /test/1.2 官方示例双版本实现
写法1:RDD 底层实现(官方经典范式)
对应官方 RDD 编程指南的标准实现,完整体现 MapReduce 执行流程。
# 1. 从 HDFS 读取文件,生成初始 RDD
lines_rdd = sc.textFile("/test/words.txt")
# 2. 拆分单词 → 映射为(单词,1)键值对 → 按单词聚合统计
word_count_rdd = lines_rdd.flatMap(lambda line: line.split(" ")) \
.map(lambda word: (word, 1)) \
.reduceByKey(lambda a, b: a + b)
# 3. 按词频倒序,取 Top5
top5 = word_count_rdd.map(lambda x: (x[1], x[0])) \
.sortByKey(ascending=False) \
.take(5)
# 4. 打印结果
print("词频 Top5:")
for cnt, word in top5:
print(f"{word}: {cnt}")写法2:DataFrame 高阶实现(官方推荐写法)
对应官方 Quick Start 中的 DataFrame + explode 实现,内置 Catalyst 优化器,性能更优,是当前企业开发主流。
from pyspark.sql import functions as sf
# 1. 读取文本生成 DataFrame
text_df = spark.read.text("/test/words.txt")
# 2. 拆分炸裂单词 → 分组统计 → 排序
word_count_df = text_df.select(
sf.explode(sf.split(text_df.value, "\s+")).alias("word")
).groupBy("word").count().orderBy("count", ascending=False)
# 3. 显示 Top5 结果
word_count_df.show(5)1.3 运行方式
交互式调试:启动
pyspark --master yarn后逐行粘贴代码,查看每一步输出脚本提交:保存为
wordcount.py,执行正式提交:spark-submit --master yarn --deploy-mode client wordcount.py
二、数值计算:Pi 圆周率计算(官方经典示例)
2.1 示例说明
官方 examples 目录中的 pi.py 是最经典的数值计算示例,使用蒙特卡洛随机投点法计算圆周率 π,无需外部数据,属于纯计算型示例,适合验证集群并行计算能力。
算法原理:在边长为 2 的正方形内画内切圆(半径为 1),随机向正方形内投点,落在圆内的概率 = 圆面积 / 正方形面积 = π/4,因此 π = 4 × 圆内点数 / 总投点数。
2.2 PySpark 实现(对齐官方逻辑)
import random
# 投点总数:数值越大精度越高,集群并行计算优势越明显
NUM_SAMPLES = 1000000
def inside(_):
"""判断随机点是否落在单位圆内"""
x, y = random.random(), random.random()
return x * x + y * y < 1
# 并行生成数据 → 过滤圆内点 → 统计数量
inside_count = sc.parallelize(range(NUM_SAMPLES)) \
.filter(inside) \
.count()
# 计算圆周率
pi = 4.0 * inside_count / NUM_SAMPLES
print(f"蒙特卡洛法计算圆周率:{pi:.6f}")2.3 直接运行官方源文件
EMR 集群已预装 Spark 官方示例,可直接运行验证环境:
# 运行官方 Python 版 Pi 示例
spark-submit --master yarn /opt/apps/SPARK3/spark-current/examples/src/main/python/pi.py三、结构化测试数据 + Spark SQL 官方示例
3.1 生成结构化测试数据
生成 CSV 格式的学生成绩数据,用于学习 DataFrame 操作与 Spark SQL 查询,对应官方 SQL 编程指南的基础案例。
# 生成学生成绩 CSV 文件
cat > student_score.csv << EOF
id,name,gender,age,subject,score
1,张三,男,20,数学,85
2,李四,女,21,数学,92
3,王五,男,19,数学,78
4,赵六,女,20,数学,88
5,孙七,男,21,语文,90
6,周八,女,19,语文,87
7,吴九,男,20,语文,76
8,郑十,女,21,语文,95
1,张三,男,20,语文,82
2,李四,女,21,语文,89
EOF
# 上传到 HDFS
hdfs dfs -put student_score.csv /test/3.2 DataFrame + SQL 双模式实现
from pyspark.sql import functions as sf
# 1. 读取 CSV 文件,自动推断表头与数据类型
score_df = spark.read.csv(
"/test/student_score.csv",
header=True,
inferSchema=True
)
# 2. 查看表结构与样本数据
score_df.printSchema()
score_df.show()
# 3. DSL 方式:统计各科目平均分、人数
subject_stats = score_df.groupBy("subject").agg(
sf.avg("score").alias("avg_score"),
sf.count("*").alias("student_cnt")
).orderBy("avg_score", ascending=False)
print("各科目统计:")
subject_stats.show()
# 4. SQL 方式:查询男生数学成绩排名
score_df.createOrReplaceTempView("v_student_score")
print("男生数学成绩排名:")
spark.sql("""
select name, score
from v_student_score
where gender = '男' and subject = '数学'
order by score desc
""").show()四、图数据 + PageRank 网页排名算法
4.1 生成图链接数据
生成网页链接关系数据,每行格式为「源页面 目标页面」,对应官方 PageRank 示例的输入格式。
# 生成网页链接关系数据
cat > links.txt << EOF
A B
A C
B C
C A
D A
D C
E B
E C
EOF
# 上传 HDFS
hdfs dfs -put links.txt /test/4.2 PageRank 简化实现(对齐官方逻辑)
PageRank 是谷歌经典的网页排名算法,也是 Spark 迭代计算的典型示例,对应官方 examples 中的 pagerank.py。
# 1. 读取链接,构建邻接表:(页面, [邻居列表])
links_rdd = sc.textFile("/test/links.txt") \
.map(lambda line: line.split(" ")) \
.map(lambda x: (x[0], x[1])) \
.distinct() \
.groupByKey()
# 2. 初始化每个页面的排名权重为 1.0
ranks_rdd = links_rdd.mapValues(lambda _: 1.0)
# 3. 迭代计算 10 次
for i in range(10):
# 计算每个页面对邻居的贡献值
contribs = links_rdd.join(ranks_rdd).flatMap(lambda item: [
(neighbor, item[1][1] / len(item[1][0]))
for neighbor in item[1][0]
])
# 更新排名:阻尼系数 0.85,防止排名收敛到0
ranks_rdd = contribs.reduceByKey(lambda a, b: a + b) \
.mapValues(lambda v: 0.15 + 0.85 * v)
# 4. 输出最终排名结果
print("PageRank 最终排名:")
for page, rank in ranks_rdd.sortBy(lambda x: -x[1]).collect():
print(f"{page}: {rank:.4f}")4.3 直接运行官方源文件
spark-submit --master yarn /opt/apps/SPARK3/spark-current/examples/src/main/python/pagerank.py五、官方示例运行与学习指南
5.1 EMR 内置官方示例位置
EMR 集群预装的 Spark 官方 Python 示例路径:
/opt/apps/SPARK3/spark-current/examples/src/main/python/包含:pi.py、wordcount.py、pagerank.py、kmeans.py、logistic_regression.py 等完整官方示例,可直接运行学习。
5.2 运行方式总结
| 运行方式 | 适用场景 | 命令示例 |
|---|---|---|
| pyspark 交互式 | 入门调试、逐行理解执行逻辑 | 启动 pyspark --master yarn 后逐行运行代码 |
| spark-submit 脚本 | 正式运行、性能测试、批量任务 | spark-submit --master yarn 脚本文件名.py |
| 直接运行官方示例 | 环境验证、对比学习 | spark-submit --master yarn /opt/apps/SPARK3/.../pi.py |
5.3 学习顺序建议
- 入门基础:先跑通 WordCount、Pi,理解 RDD 惰性求值、转换算子/行动算子、分布式计算的核心概念
- 主流开发:重点学习 DataFrame 与 Spark SQL,掌握结构化数据处理,这是当前企业开发的主流方式
- 进阶算法:最后学习 PageRank、机器学习示例,理解迭代算法、Shuffle 机制、分布式计算的优势
- 官方对照:每实现一个示例,对照官方源码和文档,理解标准写法与设计思想