PySpark 中的 RDD 操作

使用 PySpark 的 Big Data 基礎

Upendra Devisetty

Science Analyst, CyVerse

PySpark 操作總覽

  • 轉換(Transformations)會建立新的 RDD
  • 動作(Actions)會對 RDD 執行計算
使用 PySpark 的 Big Data 基礎

RDD 轉換(Transformations)

  • 轉換採用延遲求值(Lazy evaluation)

  • 基本 RDD 轉換

    • map(), filter(), flatMap(), 與 union()
使用 PySpark 的 Big Data 基礎

map() 轉換

  • map() 轉換會將函式套用到 RDD 的所有元素

map

RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)
使用 PySpark 的 Big Data 基礎

filter() 轉換

  • filter 轉換會回傳僅包含通過條件的元素之新 RDD

filter

RDD = sc.parallelize([1,2,3,4])
RDD_filter = RDD.filter(lambda x: x > 2)
使用 PySpark 的 Big Data 基礎

flatMap() 轉換

  • flatMap() 轉換會為原始 RDD 的每個元素回傳多個值

RDD = sc.parallelize(["hello world", "how are you"])
RDD_flatmap = RDD.flatMap(lambda x: x.split(" "))
使用 PySpark 的 Big Data 基礎

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 的 Big Data 基礎

RDD 動作(Actions)

  • 這類操作會在對 RDD 執行計算後回傳一個值

  • 基本 RDD 動作

    • collect()

    • take(N)

    • first()

    • count()

使用 PySpark 的 Big Data 基礎

collect() 與 take() 動作

  • collect() 會將資料集所有元素回傳為陣列

  • take(N) 會回傳含前 N 個元素的陣列

RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
使用 PySpark 的 Big Data 基礎

first() 與 count() 動作

  • first() 會輸出 RDD 的第一個元素
RDD_map.first()
[1]
  • count() 會回傳 RDD 的元素數量
RDD_flatmap.count()
5
使用 PySpark 的 Big Data 基礎

一起來練習 RDD 操作

使用 PySpark 的 Big Data 基礎

Preparing Video For Download...