在大数据开发领域中,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.

核心特点

  1. RDD 只读不可修改,只能通过算子生成新 RDD
  2. RDD 之间存在血缘依赖关系,容错靠血缘重算
  3. 惰性执行:只有触发行动算子才会真正执行任务
  4. 天然支持分区,实现集群并行运算

2. RDD 五大核心特性

  1. RDD 由多个分区组成
  2. RDD 转换 = 分区并行转换
  3. 血缘依赖机制(容错核心)
  4. KV 类型 RDD 可选自定义分区器
  5. 最优计算位置(本地优先调度)

二、怎么创建 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处理,没有返回值

七、学习总结与避坑指南

  1. 牢记转换懒执行,行动触发任务核心机制
  2. 业务优先使用 reduceByKey,少用 groupByKey 避免数据倾斜
  3. 频繁复用 RDD 一定要做缓存,减少重复计算
  4. 大集群慎用 collect,容易造成 Driver 内存溢出
  5. 算子优先选择无 Shuffle 类型,提升运行速度

更多推荐