combime 2 KeyValueGroupedDatasets

Viewed 80

I have 2 datasets which I group by key

val stgDS = Seq(("1", "1"), ("1", "2"), ("1", "3"), ("1", "4"), ("1", "5"), ("2", "1"), ("2", "2"), ("2", "3"), ("2", "4"), ("2", "5"))
      .toDF("number", "time")
      .as[Stg]

val aggDS = Seq(("1", "1"), ("1", "4"), ("1", "8"), ("2", "2"), ("2", "5"))
  .toDF("number", "time")
  .as[Agg]

After that I can apply a function to each value like so

stgDS.groupByKey(_.number)
  .flatMapGroups{case(k, iterator) => somefunction(iterator)}

How can I combine

stgDS.groupByKey(_.number)
aggDS.groupByKey(_.number)

to get something like

(k, (iteratorStg, iteratorAgg))

to then perform

.flatMapGroups{case(k, (iteratorStg, iteratorAgg)) => somefunction(iteratorStg, iteratorAgg)}

I'm looking at combineByKey finction, but either it's just another variant of grouping or I don't get how it works.

A simple join wouldn't do, because I want to loop over these iterators separately.

1 Answers

Two KeyValueGroupedDataset can be combined with cogroup.

From the cogroup doc:

Applies the given function to each cogrouped data. For each unique group, the function will be passed the grouping key and 2 iterators containing all elements in the group from Dataset this and other.

The code

val stgGroupedDS = stgDS.groupByKey(_.number)
val aggGroupedDS = aggDS.groupByKey(_.number)

stgGroupedDS.cogroup(aggGroupedDS)(
    //run whatever logic is required and then return an iterator
    (key:String, it1:Iterator[Stg], it2:Iterator[Avg])  
        => Seq((key, it1.toList.mkString(",") + "//" + it2.toList.mkString(",")))
              .iterator
)
.show(false)

prints

+---+------------------------------------------------------------------------+
|_1 |_2                                                                      |
+---+------------------------------------------------------------------------+
|1  |Stg(1,1),Stg(1,2),Stg(1,3),Stg(1,4),Stg(1,5)//Avg(1,1),Avg(1,4),Avg(1,8)|
|2  |Stg(2,1),Stg(2,2),Stg(2,3),Stg(2,4),Stg(2,5)//Avg(2,2),Avg(2,5)         |
+---+------------------------------------------------------------------------+
Related