spark-user mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From Sean Owen <>
Subject Re: GroupByKey results in OOM - Any other alternative
Date Sat, 14 Jun 2014 18:07:39 GMT
Grouping by key is always problematic since a key might have a huge number
of values. You can do a little better than grouping *all* values and *then*
finding distinct values by using foldByKey, putting values into a Set. At
least you end up with only distinct values in memory. (You don't need two
maps either, right?)

If the number of distinct values is still huge for some keys, consider the
experimental method countApproxDistinctByKey:

This should be much more performant at the cost of some accuracy.

On Sat, Jun 14, 2014 at 1:58 PM, Vivek YS <> wrote:

> Hi,
>    For last couple of days I have been trying hard to get around this
> problem. Please share any insights on solving this problem.
> Problem :
> There is a huge list of (key, value) pairs. I want to transform this to
> (key, distinct values) and then eventually to (key, distinct values count)
> On small dataset
> groupByKey().map( x => (x_1, x._2.distinct)) => (x_1,
> x._2.distinct.count))
> On large data set I am getting OOM.
> Is there a way to represent Seq of values from groupByKey as RDD and then
> perform distinct over it ?
> Thanks
> Vivek

View raw message