How to assign unique contiguous numbers to elements in a Spark RDD
apache-spark, apache-spark-mllib
Solution
Starting with Spark 1.0 there are two methods you can use to solve this easily:
- `RDD.zipWithIndex` is just like `Seq.zipWithIndex`, it adds contiguous (`Long`) numbers. This needs to count the elements in each partition first, so your input will be evaluated twice. Cache your input RDD if you want to use this.
- `RDD.zipWithUniqueId` also gives you unique `Long` IDs, but they are not guaranteed to be contiguous. (They will only be contiguous if each partition has the same number of elements.) The upside is that this does not need to know anything about the input, so it will not cause double-evaluation.
Problem
I have a dataset of `(user, product, review)`, and want to feed it into mllib's ALS algorithm. The algorithm needs users and products to be numbers, while mine are String usernames and String SKUs. Right now, I get the distinct users and SKUs, then assign numeric IDs to them outside of Spark. I was wondering whether there was a better way of doing this. The one approach I've thought of is to write a custom RDD that essentially enumerates 1 through `n`, then call zip on the two RDDs.