Spark核心知识点详解:RDD、分区与常用算子(附实战代码)
Spark是大数据领域主流分布式计算框架,核心抽象为弹性分布式数据集(RDD),分区机制与算子操作是其高效计算的关键。本文精简梳理RDD特性、分区规则、高阶函数及常用算子,搭配实战代码,帮助初学者快速掌握核心用法、避开常见坑点。
一、RDD的诞生:为什么需要弹性分布式数据集?
传统集合(如Java ArrayList、Python List)仅能在单台服务器内存中存储,无法实现分布式存储与计算,难以满足大数据处理需求。
示例代码(读取外部文件并处理):
# step1: 读取数据
input = sc.textFile("输入路径")
# step2: 处理数据(map、filter等操作)
result = input.map(lambda x: x.strip())
# step3: 保存结果
result.saveAsTextFile("输出路径")
Spark设计的弹性分布式数据集(RDD),可像集合一样存储元素,同时具备分布式特性,支持跨节点并行计算,解决了传统集合的局限。
1.1 RDD的定义与创建方式
官方定义:RDD是容错、可并行操作的分布式元素集合,核心创建方式有两种:
-
并行化集合:通过parallelize方法将Driver端集合转为RDD,指定分区数实现分布式存储。
-
读取外部存储:通过textFile、wholeTextFile等方法,读取HDFS、本地文件等外部数据转为RDD。
实战代码示例:
import os
from pyspark import SparkContext, SparkConf
if __name__ == '__main__':
# 配置环境变量(按自身环境修改路径)
os.environ['JAVA_HOME'] = 'C:\Program Files\Java\jdk1.8.0_241'
os.environ['HADOOP_HOME'] = 'D:\hadoop-3.3.1'
os.environ['PYSPARK_PYTHON'] = 'C:/ProgramData/Miniconda3/python.exe'
# 初始化SparkContext
conf = SparkConf().setMaster("local[*]").setAppName("RDD_Create_Demo")
sc = SparkContext(conf=conf)
# 方式1:并行化集合创建RDD
data = [1,2,3,4,5,6,7,8,9,10]
list_rdd = sc.parallelize(data, numSlices=2) # numSlices指定分区数
list_rdd.foreach(lambda x: print(x))
# 方式2:读取外部文件创建RDD
file_rdd = sc.textFile("../datas/filter.txt", minPartitions=2) # 最小分区数
file_rdd.foreach(lambda line: print(line))
sc.stop() # 关闭SparkContext
1.2 RDD的五大核心特性
RDD的五大特性是其容错性和并行性的基础,也是理解Spark计算模型的关键:
-
分区特性:由多个分区组成,每个分区存储在不同节点,是分布式计算的最小单元。
-
并行转换特性:转换操作(如map)对所有分区并行执行,每个分区对应一个Task。
-
依赖关系(血脉机制):保存与其他RDD的依赖关系,数据丢失可通过依赖重新构建。
-
可选分区器:KV类型RDD可自定义分区器,默认提供HashPartition和RangePartition。
-
本地优先计算:优先将Task分配到数据所在节点,减少数据传输,提升效率。
二、RDD分区的设定规则:影响计算效率的关键
分区数直接影响并行度和计算效率,不同创建方式的分区规则如下:
2.1 parallelize创建RDD的分区规则
-
指定分区数:numSlices参数直接决定分区数。
-
未指定分区数:由spark.default.parallelism参数决定,不同运行模式默认值不同。
2.2 textFile读取外部数据的分区规则
-
未指定分区数:取spark.default.parallelism和2的最小值。
-
指定分区数:为最小分区数,实际分区数按HDFS分片规则(默认128M/片)计算。
2.3 不同运行模式的默认并行度
-
Local模式:等于本地CPU核数(local[*]自动适配)。
-
Mesos模式:默认并行度为8。
-
集群模式(Standalone、Yarn):取所有Executor总核数与2的最大值。
注:子RDD默认分区数与父RDD一致,可通过repartition、coalesce算子手动修改。
三、高阶函数及Lambda表达式:Spark算子的基础
Spark算子依赖高阶函数和Lambda表达式,Python中Lambda可简化代码、提升开发效率。
3.1 Lambda表达式基础
Lambda是匿名函数,适用于函数仅使用一次、逻辑仅一行的场景,语法:lambda 参数: 表达式。
import math
list1 = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
# Lambda实现平方、立方计算
pingFang = lambda x: x * x
liFang = lambda x: math.pow(x, 3)
print(list(map(pingFang, list1))) # 平方结果
print(list(map(liFang, list1))) # 立方结果
3.2 高阶函数定义
高阶函数是参数包含另一个函数的函数,Spark的map、filter等算子均为高阶函数。
关键区别:函数本地计算,无法并行;算子分布式计算,可并行处理RDD分区数据。
四、RDD常用基础算子:实战必备
Spark算子分为两类:转换算子(Lazy加载,返回RDD)和触发算子(触发Job,返回非RDD)。
4.1 算子分类速查表
转换算子:map、filter、flatMap、reduceByKey、groupByKey、sortBy、sortByKey、union、distinct等。
触发算子:count、foreach、saveAsTextFile、first、take、collect、reduce、top等。
Shuffle算子(影响性能):reduceByKey、groupByKey、sortBy、join、distinct等。
4.2 常用转换算子实战
(1)map算子:一对一转换
对RDD每个元素执行函数,返回新RDD,适用于一对一转换。
# 计算每个元素的立方
list01 = [1,2,3,4,5,6]
listRdd = sc.parallelize(list01)
mapRdd = listRdd.map(lambda x: math.pow(x, 3))
mapRdd.foreach(lambda x: print(x))
(2)flatMap算子:扁平化转换
展开嵌套集合,类似SQL的explode函数。
# 按“/”分割每行歌曲,展开为单个歌曲名
fileRdd = sc.textFile("../datas/a.txt", 2)
flatRdd = fileRdd.flatMap(lambda line: line.split("/"))
flatRdd.foreach(lambda x: print(x))
(3)filter算子:数据过滤
按条件过滤元素,适用于行级过滤(类似SQL的where)。
import re
# 过滤状态不为-1且字段数为4的行
fileRdd = sc.textFile("../datas/b.txt", 2)
filterRdd = fileRdd.filter(lambda line: re.split(r"\s", line)[2] != '-1' and len(re.split(r"\s", line)) == 4)
filterRdd.foreach(lambda x: print(x))
(4)sortBy与sortByKey:排序算子
均实现全局排序,sortBy适用于所有RDD,sortByKey仅适用于KV类型RDD。
注:直接用foreach打印可能顺序错乱,需用collect()收集或设分区数为1。
# sortBy:按年龄降序排序
fileRdd01 = sc.textFile("../datas/c.txt")
sortByRdd = fileRdd01.sortBy(keyfunc=lambda x: int(x.split(",")[1]), ascending=False)
print(sortByRdd.collect())
# sortByKey:KV类型按Key降序排序
rdd2 = sc.parallelize([("word", 10), ("hello", 20), ("laoyan", 1)], numSlices=2)
rdd2.sortByKey(ascending=False).foreach(lambda x: print(x))
(5)groupByKey与reduceByKey:分组聚合算子
仅适用于KV类型RDD,reduceByKey先分区内聚合再Shuffle,性能优于groupByKey。
# 统计单词出现次数
rdd5 = sc.parallelize([("word", 10), ("word", 5), ("hello", 100), ("hello", 20)], numSlices=3)
# groupByKey实现(不推荐)
rdd5.groupByKey().map(lambda x: (x[0], sum(x[1]))).foreach(lambda x: print(x))
# reduceByKey实现(推荐)
rdd5.reduceByKey(lambda x, y: x + y).foreach(lambda x: print(x))
(6)分区算子:mapPartitions与foreachPartition
对整个分区操作,减少外部连接创建次数,提升性能。
# mapPartitions:计算每个分区总和
rdd1 = sc.parallelize([1,2,3,4,5,6,7,8], 2)
rdd1.mapPartitions(lambda x: [sum(x)]).foreach(print)
# foreachPartition:遍历每个分区
def showMsg(x):
for i in x:
print(i)
rdd1.foreachPartition(showMsg)
4.3 常用触发算子实战
-
count:统计元素个数;foreach:遍历无返回值;saveAsTextFile:保存数据到外部系统。
-
collect:汇总数据到Driver内存(大数据量慎用);reduce:聚合计算。
# 触发算子综合示例
rdd1 = sc.parallelize([1,2,3,4,5,6,7,8,9,10], 2)
print(rdd1.count()) # 统计个数
rdd1.foreach(lambda x: print(x*10)) # 遍历打印
rdd1.saveAsTextFile("../../output111") # 保存结果
print(rdd1.collect()) # 收集打印
print(rdd1.reduce(lambda tmp, item: tmp + item)) # 求和
五、常见问题与注意事项
-
sortBy/sortByKey打印错乱:用collect()收集或设分区数为1。
-
collect内存溢出:大数据量禁用,可用take(N)获取前N条数据。
-
Shuffle优化:减少Shuffle操作,优先用reduceByKey,合理设置分区数。
-
分区算子:优点是提升性能,缺点是单分区数据过大易内存溢出。
六、总结
本文精简梳理了Spark核心知识点,涵盖RDD特性、分区规则、高阶函数及常用算子,搭配实战代码,覆盖初学者入门核心需求。Spark的核心是“分布式”和“懒加载”,合理运用分区和算子可大幅提升计算效率。
建议多动手运行代码,熟悉算子用法和坑点,后续可深入学习Shuffle机制、数据倾斜等进阶内容。
觉得有帮助可点赞收藏,后续将持续更新Spark进阶知识点!
更多推荐



所有评论(0)