spark-user mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From Surendranauth Hiraman <>
Subject Re: GroupByKey results in OOM - Any other alternative
Date Sun, 15 Jun 2014 13:16:20 GMT

If the foldByKey solution doesn't work for you, my team uses
RDD.persist(DISK_ONLY) to avoid OOM errors.

It's slower, of course, and requires tuning other config parameters. It can
also be a problem if you do not have enough disk space, meaning that you
have to unpersist at the right points if you are running long flows.

For us, even though the disk writes are a performance hit, we prefer the
Spark programming model to Hadoop M/R. But we are still working on getting
this to work end to end on 100s of GB of data on our 16-node cluster.


On Sun, Jun 15, 2014 at 12:08 AM, Vivek YS <> wrote:

> Thanks for the input. I will give foldByKey a shot.
> The way I am doing is, data is partitioned hourly. So I am computing
> distinct values hourly. Then I use unionRDD to merge them and compute
> distinct on the overall data.
> > Is there a way to know which key,value pair is resulting in the OOM ?
> > Is there a way to set parallelism in the map stage so that, each worker
> will process one key at time. ?
> I didn't realise countApproxDistinctByKey is using hyperloglogplus. This
> should be interesting.
> --Vivek
> On Sat, Jun 14, 2014 at 11:37 PM, Sean Owen <> wrote:
>> 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


Accelerating Machine Learning

NEW YORK, NY 10001
O: (917) 525-2466 ext. 105
F: 646.349.4063
E: suren.hiraman@v <>

View raw message