Spark. Split RDD into batches

Viewed 2533

I have RDD, where each record is int:

[0,1,2,3,4,5,6,7,8]

All I need to do is split this RDD into batches. I.e. make another RDD where each element is fixed size list of elements:

[[0,1,2], [3,4,5], [6,7,8]]

This sounds trivial, however, I am puzzled last several days and cannot find anything except the following solution:

  1. Use ZipWithIndex to enumerate records in RDD:

    [0,1,2,3,4,5] -> [(0, 0),(1, 1),(2, 2),(3, 3),(4, 4),(5, 5)]

  2. Iterate over this RDD using map() and calculate index like index = int(index / batchSize)

    [1,2,3,4,5,6] -> [(0, 0),(0, 1),(0, 2),(1, 3),(1, 4),(1, 5)]

  3. Then group by generated index.

    [(0, [0,1,2]), (1, [3,4,5])]

This will get me what I need, however, I do not want to use group by here. It is trivial when you are using plain Map Reduce or some abstraction like Apache Crunch. But is there a way to produce similar result in Spark without using heavy group by?

2 Answers

You did not clearly explained why you need fixed-size RDDs, depending on what you are trying to accomplish there could be better solution, but to answer the question as it has been asked, I see the following options:
1) implement filters based on the number of items and batch sizes. For example, if you have 1000 items in the original RDD and want to split them into 10 batches, you will end up applying 10 filters, the first one checks if index is [0, 99], the second one [100, 199] and so one. After apply each filter, you will have one RDD. Important to note that the original RDD might be cached to prior to filtering. Pros: each resulting RDD can be processed separately and does not have to be fully allocated on one node. Cons: this approach becomes slower with the number of batches.
2) Logically similar to this, but instead of filter, you just implement a custom partitioner that returns a partition id based on the index (key) as described here: Custom partitioner for equally sized partitions. Pros: faster than filters. Cons: each partition has to be fit into one node.
3) If the order in the original RDD is not important and just need it to be roughly equally chunked, you can coalesce/repartition in, explained here https://jaceklaskowski.gitbooks.io/mastering-apache-spark/spark-rdd-partitions.html

Related