Principes de Big Data avec PySpark
Upendra Devisetty
Science Analyst, CyVerse


Transformations RDD de base
map(), filter(), flatMap(), et union()
RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)

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

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

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)
Opérations qui renvoient une valeur après exécution d'un calcul sur le RDD
Actions RDD de base
collect()
take(N)
first()
count()
collect() retourne tous les éléments de l'ensemble de données dans un tableau
take(N) retourne un tableau avec les N premiers éléments de l'ensemble de données
RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
RDD_map.first()
[1]
RDD_flatmap.count()
5
Principes de Big Data avec PySpark