Spark 学习必备:RDD 算子从入门到精通
在大数据开发领域中,Spark 凭借高速内存计算的优势,早已成为离线计算、实时数仓、数据分析的主流框架。而 RDD(弹性分布式数据集) 是 Spark 最核心、最底层的数据模型,想要学好 Spark,必先吃透 RDD 算子。
一、RDD 基础认知
1. 什么是 RDD
RDD 全称:Resilient Distributed Dataset
- Resilient 弹性:数据丢失可自动重算,具备容错机制
- Distributed 分布式:数据拆分存储在集群多个节点
- Dataset 数据集:不可变、可分区、支持并行计算的数据集合
官方英文释义:
Spark revolves around the concept of a resilient distributed dataset (RDD), which is a fault-tolerant collection of elements that can be operated on in parallel.
核心特点
- RDD 只读不可修改,只能通过算子生成新 RDD
- RDD 之间存在血缘依赖关系,容错靠血缘重算
- 惰性执行:只有触发行动算子才会真正执行任务
- 天然支持分区,实现集群并行运算
2. RDD 五大核心特性
- RDD 由多个分区组成
- RDD 转换 = 分区并行转换
- 血缘依赖机制(容错核心)
- KV 类型 RDD 可选自定义分区器
- 最优计算位置(本地优先调度)
二、怎么创建 RDD?两种超简单方式
Spark 官方给了两种创建 RDD 的方法,小白也能轻松上手~
方式 1:从已有的列表创建(并行化)
如果我们已经有一个普通的列表,想把它变成 RDD(让多台电脑处理),用parallelize方法就行!
# 假设我们有一个普通列表
data = [1, 2, 3, 4, 5]
# 用parallelize转换成RDD
rdd = sc.parallelize(data)
方式 2:从外部文件读取
# 读取本地或HDFS上的文本文件
rdd = sc.textFile("文件路径")
三、RDD 的 “分区”:数据怎么分?🧩
RDD 会把数据分成好几份(叫 “分区”),每一份存在不同的电脑上。分区数怎么定呢?分两种情况:
情况 1:用 parallelize 创建时
可以自己指定分区数,比如:
rdd = sc.parallelize([1,2,3,4], numSlices=2) # 分成2个分区
如果不指定,Spark 会根据你的电脑核数或集群情况自动定
情况 2:用 textFile 读取时
默认按文件大小分,一般 128M 一个分区(和 Hadoop 的块大小类似)。也可以手动指定:
rdd = sc.textFile("文件路径", minPartitions=3) # 至少3个分区
默认分区数怎么算?
Spark 有个默认参数spark.default.parallelism,自动算分区数:
- 本地模式:看你电脑的 CPU 核数
- Mesos 模式:默认 8 个
- 集群模式:所有电脑的核数总和,至少 2 个
四、前置基础:Lambda 表达式与高阶函数
1. 高阶函数定义
一个函数的参数是另一个函数,即为高阶函数,Spark 所有算子底层都是高阶函数。
2. Lambda 匿名函数
适用场景:函数仅使用一次、逻辑仅有一行,简化代码
# 普通函数
def square(x):
return x*x
# lambda匿名函数
lambda x:x*x
五、RDD 算子两大核心分类(重中之重)
Spark RDD 算子一共分为两大类,也是学习重中之重:
1. 转换算子(Transformation)
- 懒执行:调用后不立即执行,仅记录血缘关系
- 返回新 RDD,原 RDD 不变
- 分为窄依赖、宽依赖(存在 Shuffle)
常见转换算子举例:
(1)map:一对一转换
把 RDD 里的每个元素都做个处理,比如给每个数字乘 2:
# 创建一个RDD
rdd = sc.parallelize([1, 2, 3, 4])
# 用map给每个元素乘2(此时不执行) 等后面调用行动算子时,才会真正计算
map_rdd = rdd.map(lambda x: x * 2)
(2)flatMap:“压扁” 列表
如果元素是列表,flatMap 会把列表拆开,变成单个元素。比如处理句子拆单词:
rdd = sc.parallelize(["hello world", "spark is cool"])
// 按空格分词扁平化
flat_rdd = rdd.flatMap(lambda x: x.split(" "))
# 结果会是:["hello", "world", "spark", "is", "cool"]
(3)filter:过滤元素
留下符合条件的元素,比如保留偶数:
rdd = sc.parallelize(List(1,2,3,4,5,6))
// 只保留偶数
filter_rdd = rdd.filter(lambda x: x % 2 == 0) # 结果是[2,4,6]
其他转换算子:
(1)distinct 去重算子:自动去除 RDD 重复数据,底层存在 Shuffle
rdd = sc.parallelize(List(1,1,2,2,3,3))
disRdd = rdd.distinct() # 结果为 [1,2,3]
(2)groupByKey:按 key 分组
如果 RDD 是键值对,可以按 key 分组,把相同 key 的 value 放进一个列表:
# 创建键值对RDD
rdd = sc.parallelize([("spark",1), ("spark",2), ("hadoop",3)])
# 按key分组
group_rdd = rdd.groupByKey()
# 结果是:[("spark", [1,2]), ("hadoop", [3])]
(3)reduceByKey:按 key 聚合
# 统计每个单词出现次数
rdd = sc.parallelize([("spark",1), ("spark",1), ("hadoop",1)])
# 按key求和
reduce_rdd = rdd.reduceByKey(lambda x, y: x + y)
# 结果是:[("spark",2), ("hadoop",1)]
(4)sortBy 排序算子:按照指定字段正序 / 倒序排序
rdd = sc.parallelize([3,1,4,2])
sorted_rdd = rdd.sortBy(lambda x: x, ascending=False) # 降序:[4,3,2,1]
重分区算子:
如果觉得分区数不合适,可以调整:
repartition(n):任意调整分区数(会打乱数据, shuffle)coalesce(n):一般用来减少分区数(更高效,默认不 shuffle)
rdd = sc.parallelize([1,2,3,4,5], 2) # 初始2个分区
# 重分区为3个
re_rdd = rdd.repartition(3)# 或者使用
re_rdd=rdd.
coalesce(3,shuffle=True)
# 减少到1个分区
co_rdd = rdd.coalesce(1)
2. 行动算子(Action)
- 触发任务执行,提交 Job 任务
- 不返回 RDD,返回数组、数值、集合等结果
常见行动算子举例:
(1)count:统计元素个数
rdd = sc.parallelize([1,2,3,4])
print(rdd.count()) # 输出:4(触发执行)
(2)collect:收集所有元素到本地
rdd = sc.parallelize([1,2,3,4])
print(rdd.collect()) # 输出:[1,2,3,4]
数据量大时别用 collect,会撑爆本地内存!
(3)saveAsTextFile:保存结果到文件
rdd = sc.parallelize([1,2,3,4])
rdd.saveAsTextFile("输出路径")
(4)take:取前几个元素
rdd = sc.parallelize([1,2,3,4,5])
print(rdd.take(3)) # 输出:[1,2,3](前3个元素)
(4)foreach:打印
rdd = sc.parallelize([1,2,3,4,5])
rdd.foreach(print) # 输出:[1,2,3,4,5]
六、其他常用算子
1、join方面的算子 (转换算子)
两个 KV RDD 按照 key 内外连接,类似 SQL 联表查询
- join:内连接
- leftOuterJoin:左外连接
- rightOuterJoin:右外连接
rdd_singer_age = sc.parallelize([("周杰伦", 43), ("陈奕迅", 47), ("蔡依林", 41), ("林子祥", 74), ("陈升", 63)], numSlices= 2)
rdd_singer_music = sc.parallelize([("周杰伦", "青花瓷"), ("陈奕迅", "孤勇者"), ("蔡依林", "日不落"), ("林子祥", "男儿当自强"), ("动力火车", "当")], numSlices=2)# 内连接 只要相同key的 rdd1=rdd_singer_age.join(rdd_singer_music) rdd1.foreach(print) # 左外连接 要左边所有 右边只要对照上的 rdd2=rdd_singer_age.leftOuterJoin(rdd_singer_music) rdd2.foreach(print) # 右外连接 要右边所有 左边只要对照上的 rdd3=rdd_singer_age.rightOuterJoin(rdd_singer_music) rdd3.foreach(print) # 全外连接 要两边所有 rdd4=rdd_singer_age.fullOuterJoin(rdd_singer_music) rdd4.foreach(print)join展示结果:
('陈奕迅', (47, '孤勇者'))
('周杰伦', (43, '青花瓷'))
('蔡依林', (41, '日不落'))
('林子祥', (74, '男儿当自强'))
********left join 显示结果*******************************************************
('周杰伦', (43, '青花瓷'))
('蔡依林', (41, '日不落'))
('陈升', (63, None))
('陈奕迅', (47, '孤勇者'))
('林子祥', (74, '男儿当自强'))
*********right join 显示结果************************************
('动力火车', (None, '当'))
('周杰伦', (43, '青花瓷'))
('蔡依林', (41, '日不落'))
('林子祥', (74, '男儿当自强'))
('陈奕迅', (47, '孤勇者'))
********full join 显示结果*********************************************
('动力火车', (None, '当'))
('周杰伦', (43, '青花瓷'))
('蔡依林', (41, '日不落'))
('陈升', (63, None))
('陈奕迅', (47, '孤勇者'))
('林子祥', (74, '男儿当自强'))
2、分区算子
(1)mapPartitions(转换算子)
对RDD每个分区的数据进行操作,将每个分区的数据进行map转换,将转换的结果放入新的RDD中
(2)foreachParition(触发算子)
对RDD每个分区的数据进行操作,将整个分区的数据加载到内存进行foreach处理,没有返回值
七、学习总结与避坑指南
-
牢记转换懒执行,行动触发任务核心机制
-
业务优先使用 reduceByKey,少用 groupByKey 避免数据倾斜
-
频繁复用 RDD 一定要做缓存,减少重复计算
-
大集群慎用 collect,容易造成 Driver 内存溢出
-
算子优先选择无 Shuffle 类型,提升运行速度
更多推荐
所有评论(0)