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 无缝转换。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐