Các thao tác RDD trong PySpark

Nền tảng Big Data với PySpark

Upendra Devisetty

Science Analyst, CyVerse

Tổng quan thao tác PySpark

  • Transformations tạo RDD mới
  • Actions thực thi tính toán trên RDD
Nền tảng Big Data với PySpark

Các biến đổi RDD

  • Transformations theo đánh giá lười (Lazy evaluation)

  • Các biến đổi RDD cơ bản

    • map(), filter(), flatMap(), và union()
Nền tảng Big Data với PySpark

Biến đổi map()

  • Biến đổi map() áp dụng hàm cho mọi phần tử trong RDD

map

RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)
Nền tảng Big Data với PySpark

Biến đổi filter()

  • Biến đổi filter trả về RDD mới chỉ gồm các phần tử thỏa điều kiện

filter

RDD = sc.parallelize([1,2,3,4])
RDD_filter = RDD.filter(lambda x: x > 2)
Nền tảng Big Data với PySpark

Biến đổi flatMap()

  • Biến đổi flatMap() trả về nhiều giá trị cho mỗi phần tử gốc

RDD = sc.parallelize(["hello world", "how are you"])
RDD_flatmap = RDD.flatMap(lambda x: x.split(" "))
Nền tảng Big Data với PySpark

Biến đổi 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)
Nền tảng Big Data với PySpark

RDD Actions

  • Các thao tác trả về giá trị sau khi tính toán trên RDD

  • Các RDD Action cơ bản

    • collect()

    • take(N)

    • first()

    • count()

Nền tảng Big Data với PySpark

Action collect() và take()

  • collect() trả về toàn bộ phần tử của tập dữ liệu dưới dạng mảng

  • take(N) trả về mảng gồm N phần tử đầu tiên

RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
Nền tảng Big Data với PySpark

Action first() và count()

  • first() in ra phần tử đầu tiên của RDD
RDD_map.first()
[1]
  • count() trả về số phần tử trong RDD
RDD_flatmap.count()
5
Nền tảng Big Data với PySpark

Thực hành thao tác RDD

Nền tảng Big Data với PySpark

Preparing Video For Download...