Apache Spark 3.0 分布式计算:RDD 与 DataFrame

1. RDD(弹性分布式数据集)

核心概念
RDD 是 Spark 最基础的分布式数据抽象,代表不可变、分区的数据集合,支持容错和并行操作。

  • 五大特性

    1. 分区列表(Partitions)
    2. 分区计算函数(Compute Function)
    3. 依赖关系(Dependencies)
    4. 分区器(Partitioner,可选)
    5. 优先位置(Preferred Locations,可选)
  • 操作类型

    • 转换(Transformations):延迟执行,生成新 RDD(如 map, filter
    • 动作(Actions):触发计算并返回结果(如 count, collect

代码示例(Python)

# 创建 RDD
rdd = sc.parallelize([1, 2, 3, 4, 5])

# 转换:平方运算
squared_rdd = rdd.map(lambda x: x ** 2)

# 动作:求和
result = squared_rdd.reduce(lambda a, b: a + b)  # 输出:55


2. DataFrame

核心概念
DataFrame 是 Spark 1.3 引入的结构化数据抽象,以命名列(Named Columns)组织数据,类似关系型数据库表。

  • 优势
    • 内置优化器(Catalyst)自动优化执行计划
    • 支持 SQL 查询
    • 高效序列化(Tungsten引擎)
  • 数据源
    支持 JSON、Parquet、Hive 表等结构化/半结构化数据。

代码示例(Python)

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("example").getOrCreate()

# 创建 DataFrame
data = [("Alice", 34), ("Bob", 45)]
df = spark.createDataFrame(data, ["Name", "Age"])

# SQL 查询
df.createOrReplaceTempView("people")
result = spark.sql("SELECT Name FROM people WHERE Age > 40")
result.show()  # 输出:|Name| → |Bob|


3. RDD vs DataFrame 关键区别
特性RDDDataFrame
数据格式非结构化(任意对象)结构化(Schema 约束)
优化机制无自动优化Catalyst 优化器 + Tungsten 引擎
执行效率较低(Java 序列化开销)高(二进制编码 + 向量化计算)
API 类型函数式编程(Scala/Python)声明式(SQL + DSL)
使用场景细粒度控制(如图计算)结构化数据分析(如聚合、过滤)

4. Spark 3.0 的改进
  • 自适应查询执行(AQE)
    动态优化执行计划(如合并小分区),显著提升 DataFrame 性能。
  • 加速器感知调度
    支持 GPU 资源调度,加速机器学习任务。
  • 兼容性
    保留 RDD API,但推荐优先使用 DataFrame/Dataset(性能提升 2-10 倍)。

5. 如何选择?
  • 用 RDD 当
    • 需要精细控制分区逻辑
    • 处理非结构化数据(如文本流)
    • 实现自定义分布式算法
  • 用 DataFrame 当
    • 处理结构化/半结构化数据
    • 需要 SQL 接口或高效聚合
    • 追求执行性能(尤其 Spark 3.0+)

最佳实践
优先使用 DataFrame,必要时通过 df.rdd 转换为 RDD 实现复杂逻辑。

更多推荐