RDD Transformation Functions
<RDD Variable>.<transformation function>(func)map(func)
val rdd=sc.parallelize(List(1,2,3,4,5))
rdd.map(x=>x+1).collect.foreach(println)
/*
2
3
4
5
6
*/filter(func)
val rdd=sc.parallelize(List(1,2,3,4,5))
rdd.filter(x=>x%2==0).collect.foreach(println)
/*
2
4
*/flatMap(func)
mapPartitions(func)
mapPartitionsWithIndex(func)
sample(withReplacement, fraction, seed)
union(otherDataset)
intersection(otherDataset)
distinct([numPartitions]))
groupByKey([numPartitions])
reduceByKey(func, [numPartitions])
aggregateByKey(zeroValue)(seqOp, combOp, [numPartitions])
join(otherDataset, [numPartitions])
cogroup(otherDataset, [numPartitions])
cartesian(otherDataset)
pipe(command, [envVars])
coalesce(numPartitions)
repartition(numPartitions)
repartitionAndSortWithinPartitions(partitioner)
Last updated