THE SCENARIO
I'm trying to write a Spark program that efficiently performs a left outer join between two RDDs. One caveat is that these RDDs can have duplicate keys, which apparently causes the whole program to be inefficient.
What I'm trying to achieve is simple:
- Given two RDDs:
rdd1andrdd2(both have the same structure:(k, v)) - Using
rdd1andrdd2, generate another RDDrdd3that has the structure:(k1, v1, List(v2..)) k1andv1come fromrdd1(same values, this will lead tordd1andrdd3have the same length)List(v2..)is a list whose values are coming from the values ofrdd2- To add an
rdd2'svto the list inrdd3's tuple, itsk(the key fromrdd2) should match thekfromrdd1
MY ATTEMPT
My approach was to use a left outer join. So, I came up with something like this:
rdd1.leftOuterJoin(rdd2).map{case(k, (v1, v2)) => ((k, v1), Array(v2))}
.reduceByKey(_ ++ _)
This actually produces the result that I'm trying to acheive. But, when I use a huge data, the program becomes very slow.
AN EXAMPLE
Just in case my idea is not clear yet, I have the following example:
Given two RDDs that have the following data:
rdd1:
key | value
-----------
1 | a
1 | b
1 | c
2 | a
2 | b
3 | c
rdd2:
key | value
-----------
1 | v
1 | w
1 | x
1 | y
1 | z
2 | v
2 | w
2 | x
3 | y
4 | z
The resulting rdd3 should be
key | value | list
------------------------
1 | a | v,w,x,y,z
1 | b | v,w,x,y,z
1 | c | v,w,x,y,z
2 | a | v,w,x
2 | b | v,w,x
3 | c | y