Apache Spark 3.0 分布式计算:RDD 与 DataFrame
·
Apache Spark 3.0 分布式计算:RDD 与 DataFrame
1. RDD(弹性分布式数据集)
-
核心概念:
分布式内存中的不可变数据集合,支持并行操作。
基本特性:- 分区存储(Partitioned)
- 容错性(通过血缘关系 Lineage 重建)
- 惰性求值(Lazy Evaluation)
-
操作示例:
from pyspark import SparkContext sc = SparkContext("local", "RDD Example") # 创建RDD data = [1, 2, 3, 4, 5] rdd = sc.parallelize(data, numSlices=3) # 分3个分区 # 转换操作:map squared_rdd = rdd.map(lambda x: x**2) # 行动操作:collect print(squared_rdd.collect()) # 输出: [1, 4, 9, 16, 25]
2. DataFrame
-
核心概念:
结构化数据抽象,基于命名列的分布式数据集。
关键优势:- Catalyst 优化器自动优化执行计划
- Tungsten 引擎提升内存/CPU效率
- 支持 SQL 查询和 Schema 推断
-
操作示例:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("DataFrame Example").getOrCreate() # 创建DataFrame data = [("Alice", 34), ("Bob", 45)] columns = ["Name", "Age"] df = spark.createDataFrame(data, columns) # SQL式查询 df.filter("Age > 40").show() # 输出: # +-----+---+ # | Name|Age| # +-----+---+ # | Bob| 45| # +-----+---+
3. 核心区别对比
| 特性 | RDD | DataFrame |
|---|---|---|
| 数据格式 | 非结构化对象 | 结构化命名列 |
| 优化机制 | 手动优化 | Catalyst 自动优化 |
| 执行效率 | 较低(JVM对象开销) | 较高(Tungsten二进制格式) |
| API 类型 | 函数式编程(map/reduce) | 声明式(SQL/DSL) |
| 类型安全 | 编译时检查 | 运行时检查 |
| 适用场景 | 细粒度控制、复杂算法 | 结构化分析、ETL |
4. Spark 3.0 关键优化
-
自适应查询(AQE):
动态合并 shuffle 分区,优化资源利用。
例如:初始分区数 $n$ 根据数据量自动调整为 $k$($k < n$)。 -
动态分区修剪(DPP):
跳过无关数据分区的扫描,减少 I/O。
满足条件:$ \text{filter}(pred) \implies \text{skip_partition}(p) $ -
加速器支持:
集成 GPU 加速计算(需配合 Rapids 插件)。
5. 使用建议
-
选 RDD 当:
- 需精细控制分区逻辑
- 处理非结构化数据(如文本流)
- 实现自定义序列化
-
选 DataFrame 当:
- 执行 SQL 查询或聚合操作
- 需要自动优化执行计划
- 与外部数据源集成(Parquet/CSV)
总结:在 Spark 3.0 中,DataFrame 是首选——其优化器可提升性能 $2\text{-}10\times$,而 RDD 更适合底层定制化需求。两者可通过
df.rdd无缝转换。
更多推荐


所有评论(0)