Spark基础:RDD创建和转换
·
Spark基础:RDD创建和转换
1. RDD创建
RDD(弹性分布式数据集)是Spark的核心抽象,可通过以下方式创建:
a. 从集合创建(并行化)
使用sc.parallelize()将本地集合转为分布式RDD:
from pyspark import SparkContext
sc = SparkContext("local", "RDD Demo")
data = [1, 2, 3, 4, 5]
rdd = sc.parallelize(data) # 分区数默认为系统核心数
b. 从外部数据源创建
支持HDFS、本地文件系统等:
# 从文本文件创建(每行为一个元素)
text_rdd = sc.textFile("hdfs://path/to/file.txt")
# 从CSV创建(需预处理)
csv_rdd = text_rdd.map(lambda line: line.split(","))
2. RDD转换操作
转换操作是惰性的,仅记录计算逻辑,触发行动操作时才执行。常用转换:
a. map(func)
对每个元素应用函数,输入输出为$1:1$关系:
squared_rdd = rdd.map(lambda x: x**2) # [1,4,9,16,25]
b. filter(func)
保留满足条件的元素:
even_rdd = rdd.filter(lambda x: x % 2 == 0) # [2,4]
c. flatMap(func)
每个元素映射为0-N个输出,输出为扁平化结构:
words_rdd = text_rdd.flatMap(lambda line: line.split(" "))
d. 键值对转换
针对(key, value)格式RDD:
pair_rdd = rdd.map(lambda x: (x % 2, x)) # [(1,1), (0,2), (1,3)...]
reduce_rdd = pair_rdd.reduceByKey(lambda a, b: a + b) # 分组求和
3. 转换操作特性
- 惰性求值:转换操作仅构建DAG(有向无环图),不立即计算
- 不可变性:每次转换生成新RDD,原RDD不变
- 依赖关系:窄依赖(如
map)与宽依赖(如reduceByKey)- 窄依赖:父RDD分区到子RDD分区为$1:1$映射
- 宽依赖:涉及Shuffle操作,如分组聚合
4. 完整示例
# 创建RDD并执行转换链
rdd = sc.parallelize([10, 20, 30, 40])
result = (
rdd.map(lambda x: x/10) # -> [1.0, 2.0, 3.0, 4.0]
.filter(lambda x: x > 2) # -> [3.0, 4.0]
.flatMap(lambda x: [x, x*10]) # -> [3.0, 30.0, 4.0, 40.0]
)
# 行动操作触发计算(如collect)
print(result.collect()) # 输出: [3.0, 30.0, 4.0, 40.0]
关键点:
- 转换操作定义计算流程,行动操作(如
collect(),count())触发实际执行- 通过转换链实现高效分布式计算,避免中间结果落盘
- 宽依赖操作(如
groupByKey)需谨慎使用,可能引发Shuffle性能瓶颈
更多推荐
所有评论(0)