PySpark 中的 RDD 操作

使用 PySpark 的大数据基础

Upendra Devisetty

Science Analyst, CyVerse

PySpark 操作概览

  • 转换会创建新的 RDD
  • 动作在 RDD 上执行计算
使用 PySpark 的大数据基础

RDD 转换(Transformations)

  • 转换遵循惰性求值(Lazy evaluation)

  • 基本 RDD 转换

    • map(), filter(), flatMap(), union()
使用 PySpark 的大数据基础

map() 转换

  • map() 对 RDD 中所有元素应用函数

map

RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)
使用 PySpark 的大数据基础

filter() 转换

  • filter 转换返回满足条件的元素组成的新 RDD

filter

RDD = sc.parallelize([1,2,3,4])
RDD_filter = RDD.filter(lambda x: x > 2)
使用 PySpark 的大数据基础

flatMap() 转换

  • flatMap() 为原 RDD 的每个元素返回多个值

RDD = sc.parallelize(["hello world", "how are you"])
RDD_flatmap = RDD.flatMap(lambda x: x.split(" "))
使用 PySpark 的大数据基础

union() 转换

inputRDD = sc.textFile("logs.txt")
errorRDD = inputRDD.filter(lambda x: "error" in x.split())
warningsRDD = inputRDD.filter(lambda x: "warnings" in x.split())
combinedRDD = errorRDD.union(warningsRDD)
使用 PySpark 的大数据基础

RDD 动作(Actions)

  • 在 RDD 上运行计算后返回一个值的操作

  • 基本 RDD 动作

    • collect()

    • take(N)

    • first()

    • count()

使用 PySpark 的大数据基础

collect() 与 take() 动作

  • collect() 将数据集的所有元素作为数组返回

  • take(N) 返回包含前 N 个元素的数组

RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
使用 PySpark 的大数据基础

first() 与 count() 动作

  • first() 返回 RDD 的第一个元素
RDD_map.first()
[1]
  • count() 返回 RDD 中元素个数
RDD_flatmap.count()
5
使用 PySpark 的大数据基础

练习 RDD 操作

使用 PySpark 的大数据基础

Preparing Video For Download...