[ PySpark RDD (Resilient Distributed Datasets) Operations ] [ cheatsheet ]
1. RDD Creation
● Create RDD from a list: rdd = [Link]([1, 2, 3, 4, 5])
● Create RDD from a file: rdd = [Link]("[Link]")
● Create RDD from a directory: rdd = [Link]("directory/")
● Create empty RDD: rdd = [Link]()
2. RDD Transformations
● Map: [Link](lambda x: x * 2)
● FlatMap: [Link](lambda x: [x, x * 2, x * 3])
● Filter: [Link](lambda x: x > 10)
● Distinct: [Link]()
● Sample: [Link](withReplacement=True, fraction=0.5)
● Union: [Link](rdd2)
● Intersection: [Link](rdd2)
● Subtract: [Link](rdd2)
● Cartesian: [Link](rdd2)
● Zip: [Link](rdd2)
● ZipWithIndex: [Link]()
● GroupBy: [Link](lambda x: x % 2)
● SortBy: [Link](lambda x: x, ascending=False)
● PartitionBy: [Link](3)
● MapPartitions: [Link](lambda partition: [x * 2 for x in
partition])
● MapPartitionsWithIndex: [Link](lambda index,
partition: [(index, x) for x in partition])
● FlatMapValues: [Link](lambda x: [x, x * 2])
● CombineByKey: [Link](lambda value: (value, 1), lambda acc,
value: (acc[0] + value, acc[1] + 1), lambda acc1, acc2: (acc1[0] +
acc2[0], acc1[1] + acc2[1]))
● FoldByKey: [Link](0, lambda acc, value: acc + value)
● ReduceByKey: [Link](lambda a, b: a + b)
● AggregateByKey: [Link]((0, 0), lambda acc, value: (acc[0] +
value, acc[1] + 1), lambda acc1, acc2: (acc1[0] + acc2[0], acc1[1] +
acc2[1]))
● Join: [Link](rdd2)
By: Waleed Mousa
● LeftOuterJoin: [Link](rdd2)
● RightOuterJoin: [Link](rdd2)
● FullOuterJoin: [Link](rdd2)
● Cogroup: [Link](rdd2)
3. RDD Actions
● Collect: [Link]()
● Take: [Link](5)
● First: [Link]()
● Count: [Link]()
● CountByValue: [Link]()
● Reduce: [Link](lambda a, b: a + b)
● Aggregate: [Link]((0, 0), lambda acc, value: (acc[0] + value,
acc[1] + 1), lambda acc1, acc2: (acc1[0] + acc2[0], acc1[1] + acc2[1]))
● Fold: [Link](0, lambda acc, value: acc + value)
● Max: [Link]()
● Min: [Link]()
● Sum: [Link]()
● Mean: [Link]()
● Variance: [Link]()
● Stdev: [Link]()
● TakeSample: [Link](withReplacement=True, num=5)
● Foreach: [Link](lambda x: print(x))
● ForeachPartition: [Link](lambda partition: [print(x) for x
in partition])
● Top: [Link](5)
● TakeOrdered: [Link](5, lambda x: -x)
● SaveAsTextFile: [Link]("output/")
● SaveAsPickleFile: [Link]("output/")
4. RDD Persistence and Caching
● Cache: [Link]()
● Persist: [Link](storageLevel=[Link].MEMORY_AND_DISK)
● Unpersist: [Link]()
● Checkpoint: [Link]()
5. RDD Partitioning
By: Waleed Mousa
● Repartition: [Link](numPartitions=10)
● Coalesce: [Link](numPartitions=5)
● GetNumPartitions: [Link]()
● Glom: [Link]()
● Zip: [Link](rdd2)
● ZipPartitions: [Link](rdd2, lambda partition1, partition2: [x
+ y for x, y in zip(partition1, partition2)])
6. RDD Set Operations
● Union: [Link](rdd2)
● Intersection: [Link](rdd2)
● Subtract: [Link](rdd2)
● Cartesian: [Link](rdd2)
7. RDD Pair Operations
● ReduceByKey: [Link](lambda a, b: a + b)
● GroupByKey: [Link]()
● CombineByKey: [Link](lambda value: (value, 1), lambda acc,
value: (acc[0] + value, acc[1] + 1), lambda acc1, acc2: (acc1[0] +
acc2[0], acc1[1] + acc2[1]))
● AggregateByKey: [Link]((0, 0), lambda acc, value: (acc[0] +
value, acc[1] + 1), lambda acc1, acc2: (acc1[0] + acc2[0], acc1[1] +
acc2[1]))
● FoldByKey: [Link](0, lambda acc, value: acc + value)
● Join: [Link](rdd2)
● LeftOuterJoin: [Link](rdd2)
● RightOuterJoin: [Link](rdd2)
● FullOuterJoin: [Link](rdd2)
● MapValues: [Link](lambda x: x * 2)
● FlatMapValues: [Link](lambda x: [x, x * 2])
● CountByKey: [Link]()
● LookupByKey: [Link](key)
● SortByKey: [Link](ascending=False)
● CoGroupByKey: [Link](rdd2)
8. RDD Numeric Operations
● Sum: [Link]()
By: Waleed Mousa
● Mean: [Link]()
● Variance: [Link]()
● Stdev: [Link]()
● HistogramByKey: [Link](buckets=10)
9. RDD Data Formats
● CSV: rdd = [Link]("[Link]").map(lambda line: [Link](","))
● JSON: rdd = [Link]("[Link]").map(lambda line: [Link](line))
● Parquet: rdd = [Link]("[Link]").rdd
● Avro: rdd =
[Link]("[Link]").load("[Link]").rdd
● SequenceFile: rdd = [Link]("[Link]", keyClass, valueClass)
10. RDD Compression
● Compress: [Link](lambda x: (x, None)).saveAsSequenceFile("output/",
codec="[Link]")
● Decompress: rdd = [Link]("output/", keyClass=None,
valueClass=None, codec="[Link]")
11. RDD Serialization
● Kryo Serializer: conf = SparkConf().set("[Link]",
"[Link]")
● Java Serializer: conf = SparkConf().set("[Link]",
"[Link]")
12. RDD Partitioner
● HashPartitioner: [Link](numPartitions=10, partitionFunc=lambda
x: hash(x))
● RangePartitioner: [Link](numPartitions=10, partitionFunc=lambda
x: int(x / 10))
● CustomPartitioner: [Link](numPartitions=10, partitionFunc=lambda
x: customPartitionFunction(x))
13. RDD Execution
● Collect: [Link]()
By: Waleed Mousa
● Reduce: [Link](lambda a, b: a + b)
● Aggregate: [Link]((0, 0), lambda acc, value: (acc[0] + value,
acc[1] + 1), lambda acc1, acc2: (acc1[0] + acc2[0], acc1[1] + acc2[1]))
● Take: [Link](5)
● First: [Link]()
● TakeSample: [Link](withReplacement=True, num=5)
● Count: [Link]()
● CountByValue: [Link]()
● Foreach: [Link](lambda x: print(x))
● ForeachPartition: [Link](lambda partition: [print(x) for x
in partition])
● CollectAsMap: [Link]()
14. RDD Debugging
● Logging: [Link](lambda x: print("Processing:", x))
● Caching: [Link]()
● Checkpointing: [Link]()
● Debugging: [Link](lambda x: x == debug_value).collect()
15. RDD Optimization
● Caching: [Link]()
● Repartitioning: [Link](numPartitions=10)
● Coalesce: [Link](numPartitions=5)
● Broadcast Variables: broadcast_var = [Link](large_data)
● Accumulator Variables: accumulator = [Link](0)
● Kryo Serialization: conf = SparkConf().set("[Link]",
"[Link]")
● Spark SQL: df = [Link](rdd, schema)
● DataFrame Operations: [Link]([Link] > 18)
By: Waleed Mousa