拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Spark 入门第一课:用 RDD 实现 WordCount,理解分布式计算的“搬砖”哲学

关键词Spark、RDD、WordCount、转换与行动、惰性求值、Shuffle引言如果你刚接触 Apache Spark第一个要写的程序几乎一定是WordCount——就像学习编程语言时写“Hello World”一样。但 WordCount 的意义远不止数数单词它浓缩了 Spark 最核心的抽象RDD弹性分布式数据集。通过这个简单的例子你将直观地理解“转换是惰性的”、“行动触发计算”、“Shuffle 代价高昂”这些 Spark 底层原理。本文将在 Docker 单机 Spark 集群上运行一个完整的 WordCount 脚本并逐行解读代码、展示真实输出最后给出原理解析。读完本文你就能摸到 Spark 的“地基”。实验环境Spark 版本4.2.0镜像apache/spark:4.2.0-scala2.13-java17-python3-ubuntu部署方式Docker 容器内 Standalone 模式1 Master 1 Worker数据目录容器内/workspace/data/text/脚本位于/workspace/scripts/准备数据我们在data/text/下放两个小文本文件内容为 Spark 相关的英文句子确保有足够的重复词1.txtApache Spark is a unified analytics engine for large-scale data processing. Spark provides high-level APIs in Python, Scala, Java and R.2.txtSpark SQL, Spark Streaming, MLlib and GraphX are built on Spark Core. Spark runs on a cluster: the master schedules, workers execute.完整代码wordcount.py实验一单词计数WordCount—— RDD API 版Spark 的 Hello World 用法 spark-submit --master master-url wordcount.py 输入目录 输出目录 importreimportsysfrompysparkimportSparkContext WORDSre.compile(r\w)defsplit_words(line):把一行拆成小写单词列表顺手去掉标点returnWORDS.findall(line.lower())defmain():input_pathsys.argv[1]output_pathsys.argv[2]scSparkContext(appNameWordCount)# master 地址由 spark-submit 的 --master 指定counts(sc.textFile(input_path).flatMap(split_words)# 一行 - 多个词.map(lambdaword:(word,1))# 词 - (词, 1).reduceByKey(lambdaa,b:ab).cache()# 下面有两个 actionsave/take缓存避免重复计算)counts.saveAsTextFile(output_path)# action触发计算写出 part-* 文件print( 词频 Top 10 )forword,cntincounts.takeOrdered(10,keylambdax:-x[1]):print(f{word}:{cnt})sc.stop()if__name____main__:main()运行命令# 清除旧输出目录saveAsTextFile 不允许覆盖dockerexecsparkrm-rf/workspace/data/out/wordcount# 提交作业到 Standalone 集群dockerexecspark /opt/spark/bin/spark-submit\--masterspark://localhost:7077\/workspace/scripts/wordcount.py\/workspace/data/text /workspace/data/out/wordcount真实输出控制台打印 Top 10 词频 词频 Top 10 the: 9 spark: 7 master: 4 in: 3 tasks: 3 workers: 3 and: 3 run: 2 driver: 2 is: 2输出目录下生成了两个 part 文件因为输入有两个文件对应两个分区$lsdata/out/wordcount/ _SUCCESS part-00000 part-00001 $headpart-00000(master,4)(coordinates,1)(in,3)我们可以在宿主机上用 shell 命令独立校验结果完全一致说明集群计算正确。原理解读RDD 与“施工队”比喻RDDResilient Distributed Dataset是 Spark 最底层的抽象一个被拆成多个分区、分布在集群各节点上的不可变数据集。你可以把它想象成一张“工料单”每个分区就是一位工人手里的一沓料单。本例中两个输入文件 → 2 个分区 → 2 个输出 part 文件一一对应。程序中的操作分为两大类类型比喻本实验例子转换Transformation改图纸但不动手flatMap、map、reduceByKey行动Action喊“开工”真正干活saveAsTextFile、takeOrdered关键认知转换是惰性的。调用flatMap、map等只是构建了一个计算蓝图并不实际执行。直到遇到第一个Action如saveAsTextFileSpark 才会从源头开始执行整个流水线。本脚本有两个 Action如果不加.cache()Spark 会重新计算两次加上缓存后第一次计算的结果被保存第二次直接复用提高效率。Shuffle 的代价reduceByKey会将相同 key 的(word, 1)对汇聚到一起这需要跨分区传输数据即 Shuffle是分布式计算中最昂贵的操作之一。理解哪些操作会触发 Shuffle是 Spark 调优的第一步。容错性Resilient如果某个分区丢失Spark 可以根据血缘Lineage——记录的转换历史——从原始数据重新计算该分区这就是“弹性”的由来。常见问题输出目录已存在saveAsTextFile不会覆盖需要先rm -rf或换路径。--master local[*]与spark://...的区别前者是本地模式driver 自己开线程模拟集群看不到真正的 Master/Worker后者提交到 Standalone 集群能观察调度细节。中文文本需要分词器或按字切分不能直接用正则\w。总结WordCount 虽小但已经展示了 RDD 的“转换-行动”模型、惰性求值、Shuffle 和容错机制。这是 Spark 一切高级 APIDataFrame、SQL的基础。掌握了它你就拿到了进入 Spark 世界的门票。下一篇博客我们将从“搬砖”升级到“看图纸”——用 DataFrame API 处理多层嵌套 JSONL 数据体验结构化计算的便利。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门