Podstawy Big Data z PySpark
Upendra Devisetty
Science Analyst, CyVerse
Akcja reduce(func) służy do agregowania elementów zwykłego RDD
Funkcja powinna być przemienna (zmiana kolejności operandów nie zmienia wyniku) i łączna
Przykład akcji reduce() w PySpark
x = [1,3,4,6]
RDD = sc.parallelize(x)
RDD.reduce(lambda x, y : x + y)
14
saveAsTextFile() zapisuje RDD do pliku tekstowego w katalogu — każda partycja jako osobny plikRDD.saveAsTextFile("tempFile")
coalesce() umożliwia zapisanie RDD jako jednego pliku tekstowegoRDD.coalesce(1).saveAsTextFile("tempFile")
Akcje RDD dostępne dla par RDD w PySpark
Akcje na parach RDD wykorzystują dane w formacie klucz-wartość
Przykłady akcji na parach RDD
countByKey()
collectAsMap()
countByKey() dostępna tylko dla typu (K, V)
Akcja countByKey() zlicza elementy dla każdego klucza
Przykład countByKey() na prostej liście
rdd = sc.parallelize([("a", 1), ("b", 1), ("a", 1)])
for kee, val in rdd.countByKey().items():
print(kee, val)
('a', 2)
('b', 1)
collectAsMap() zwraca pary klucz-wartość z RDD jako słownik
Przykład collectAsMap() na prostej krotce
sc.parallelize([(1, 2), (3, 4)]).collectAsMap()
{1: 2, 3: 4}
Podstawy Big Data z PySpark