Операции с RDD в PySpark

Основы Big Data с PySpark

Upendra Devisetty

Science Analyst, CyVerse

Обзор операций PySpark

  • Трансформации создают новые RDD
  • Действия выполняют вычисления над RDD
Основы Big Data с PySpark

Трансформации RDD

  • Трансформации используют «ленивые» вычисления

  • Основные трансформации 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

Давайте потренируемся!

Основы Big Data с PySpark

Preparing Video For Download...