Operacje na RDD w PySpark

Podstawy Big Data z PySpark

Upendra Devisetty

Science Analyst, CyVerse

Przegląd operacji w PySpark

  • Transformacje tworzą nowe RDD
  • Akcje wykonują obliczenia na RDD
Podstawy Big Data z PySpark

Transformacje RDD

  • Transformacje stosują leniwą ewaluację

  • Podstawowe transformacje RDD

    • map(), filter(), flatMap() i union()
Podstawy Big Data z PySpark

Transformacja map()

  • Transformacja map() stosuje funkcję do każdego elementu RDD

map

RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)
Podstawy Big Data z PySpark

Transformacja filter()

  • Transformacja filter() zwraca nowe RDD tylko z elementami spełniającymi warunek

filter

RDD = sc.parallelize([1,2,3,4])
RDD_filter = RDD.filter(lambda x: x > 2)
Podstawy Big Data z PySpark

Transformacja flatMap()

  • Transformacja flatMap() zwraca wiele wartości dla każdego elementu oryginalnego RDD

RDD = sc.parallelize(["hello world", "how are you"])
RDD_flatmap = RDD.flatMap(lambda x: x.split(" "))
Podstawy Big Data z PySpark

Transformacja 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)
Podstawy Big Data z PySpark

Akcje RDD

  • Operacje zwracające wartość po wykonaniu obliczeń na RDD

  • Podstawowe akcje RDD

    • collect()

    • take(N)

    • first()

    • count()

Podstawy Big Data z PySpark

Akcje collect() i take()

  • collect() zwraca wszystkie elementy zbioru danych jako tablicę

  • take(N) zwraca tablicę z pierwszymi N elementami zbioru

RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
Podstawy Big Data z PySpark

Akcje first() i count()

  • first() zwraca pierwszy element RDD
RDD_map.first()
[1]
  • count() zwraca liczbę elementów w RDD
RDD_flatmap.count()
5
Podstawy Big Data z PySpark

Czas na ćwiczenia z operacjami RDD

Podstawy Big Data z PySpark

Preparing Video For Download...