Операції RDD у PySpark

Основи Big Data з PySpark

Upendra Devisetty

Science Analyst, CyVerse

Огляд операцій PySpark

  • Перетворення створюють нові RDD
  • Дії виконують обчислення над RDD
Основи Big Data з PySpark

Перетворення RDD

  • Перетворення виконуються з ледачим обчисленням (Lazy evaluation)

  • Базові перетворення RDD

    • map(), filter(), flatMap(), і union()
Основи Big Data з PySpark

Перетворення map()

  • Перетворення map() застосовує функцію до всіх елементів RDD

map

RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)
Основи Big Data з PySpark

Перетворення filter()

  • Перетворення filter повертає новий RDD лише з елементами, що проходять умову

filter

RDD = sc.parallelize([1,2,3,4])
RDD_filter = RDD.filter(lambda x: x > 2)
Основи Big Data з PySpark

Перетворення flatMap()

  • Перетворення flatMap() повертає кілька значень для кожного елемента початкового RDD

RDD = sc.parallelize(["hello world", "how are you"])
RDD_flatmap = RDD.flatMap(lambda x: x.split(" "))
Основи Big Data з 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)
Основи Big Data з PySpark

Дії RDD

  • Це операції, що повертають значення після обчислення над RDD

  • Базові дії RDD

    • collect()

    • take(N)

    • first()

    • count()

Основи Big Data з PySpark

Дії collect() і take()

  • collect() повертає всі елементи набору даних як масив

  • take(N) повертає масив із перших N елементів набору даних

RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
Основи Big Data з PySpark

Дії first() і count()

  • first() виводить перший елемент RDD
RDD_map.first()
[1]
  • count() повертає кількість елементів у RDD
RDD_flatmap.count()
5
Основи Big Data з PySpark

Давайте потренуємось з операціями RDD

Основи Big Data з PySpark

Preparing Video For Download...