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性能瓶颈

更多推荐