# Aggregation running very slow

**URL:** <https://discuss.elastic.co/t/aggregation-running-very-slow/28192>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [August 27, 2015, 3:29pm UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192 "2015-08-27T15:29:16Z")\
**Posts on this page:** 9\
**Page:** 1

<div class="post-metadata">

**Author:** ![Philip\_K\_Adetiloye](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/philip_k_adetiloye/32/14386_2.png) [@Philip\_K\_Adetiloye](https://discuss.elastic.co/u/Philip_K_Adetiloye)\
**Post date:** [August 27, 2015, 3:29pm UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/1 "2015-08-27T15:29:16Z")

</div>

I'm using sparkSQL to query ES but aggregation (GROUP BY) is running very slow. I can't remember but seems I read that aggregation queries are also push down to ES. If that's the case, why is aggregation running so slow and how can I improve that ?

---

<div class="post-metadata">

**Author:** ![costin](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/costin/32/44950_2.png) [@costin](https://discuss.elastic.co/u/costin)\
**Post date:** [August 28, 2015, 9:57am UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/2 "2015-08-28T09:57:26Z")

</div>

When in doubt, refer to the ES-Hadoop reference documentation. Spark performs group by aggregation by itself and does not expose a filter to be pushed down.  
In other words there's nothing the connector can do for this operation.

---

<div class="post-metadata">

**Author:** ![Philip\_K\_Adetiloye](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/philip_k_adetiloye/32/14386_2.png) [@Philip\_K\_Adetiloye](https://discuss.elastic.co/u/Philip_K_Adetiloye)\
**Post date:** [September 1, 2015, 1:50pm UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/3 "2015-09-01T13:50:35Z")

</div>

@costin  
I now understand the mechanism of how ES-Hadoop work, Thanks!

If I may ask, how can I get all my data indexed in Elasticsearch to Spark ?

I'm not sure if there is a better way than using spark sql to do a `select * from BigIndex`

---

<div class="post-metadata">

**Author:** ![costin](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/costin/32/44950_2.png) [@costin](https://discuss.elastic.co/u/costin)\
**Post date:** [September 1, 2015, 2:04pm UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/4 "2015-09-01T14:04:12Z")

</div>

That is as well [explained](https://www.elastic.co/guide/en/elasticsearch/hadoop/current/spark.html#spark-write) in the docs - in fact Spark is the longest chapter in the reference documentation.

---

<div class="post-metadata">

**Author:** ![sergiokozlov](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/sergiokozlov/32/4566_2.png) [@sergiokozlov](https://discuss.elastic.co/u/sergiokozlov)\
**Post date:** [September 3, 2015, 2:34pm UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/5 "2015-09-03T14:34:48Z")

</div>

@costin Can you please elaborate more on how aggregations are handled for the ES-SparkSQL integration. (Sorry I couldn't find it in the docs). For example what will happen if I will issue a following query:

```
SELECT termA, SUM(measure) from index
WHERE termB = 'XXX' and termC > 5
GROUP BY termA

```

As I understand it should happen as:

1. ES issues query `SELECT term, measure from index WHERE termB = 'XXX' and termC > 5` on every shard (using collocation where possible)
2. This data is loaded into Spark on each node and there it's aggregated.

By this we benefit from ES index - we quickly get only filtered data and also spark computation for the GROUP BY operation.

Do I understand it correctly - can you please correct me if necessary.

Thanks,  
Sergey.

---

<div class="post-metadata">

**Author:** ![costin](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/costin/32/44950_2.png) [@costin](https://discuss.elastic.co/u/costin)\
**Post date:** [September 4, 2015, 9:03am UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/6 "2015-09-04T09:03:38Z")

</div>

[https://www.elastic.co/guide/en/elasticsearch/hadoop/current/spark.html#spark-pushdown](https://www.elastic.co/guide/en/elasticsearch/hadoop/current/spark.html#spark-pushdown)

P.S. In the future, please open a new thread instead of hijacking an existing one. Thanks!

---

<div class="post-metadata">

**Author:** ![sergiokozlov](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/sergiokozlov/32/4566_2.png) [@sergiokozlov](https://discuss.elastic.co/u/sergiokozlov)\
**Post date:** [September 4, 2015, 9:22am UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/7 "2015-09-04T09:22:00Z")

</div>

@costin Thank you for the prompt response, I saw this documentation, but it doesn't say whether aggregations are supported by push-down or not (I assume that not). That's why I asked for the steps of GROUP BY query.

Thanks,  
Sergey.

---

<div class="post-metadata">

**Author:** ![costin](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/costin/32/44950_2.png) [@costin](https://discuss.elastic.co/u/costin)\
**Post date:** [September 4, 2015, 9:40am UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/8 "2015-09-04T09:40:38Z")

</div>

Point taken. I'll update the docs to make this clear.  
Note that as I've indicated on this thread - the issue is that without hooks in Spark, aggregations-like operations are not possible since the storage (in this case Elastic) is unaware of what's going on.  
Take for example - `count`. On an `RDD` level, the connector supports this operation and instead of iterating manually through the entire content, returns the count right away.  
However on a `DataFrame` the API doesn't allow `count` to be pushed down so one will actually end up iterating over the entire data. This is true for aggregations, including group-by.

The pushdown however applies at query level, in particular SQL where the `WHERE` and `SELECT` are properly exposed and thus picked up by ES-Hadoop.

Maybe in the future, more and more operations will be exposed by Spark.

---

<div class="post-metadata">

**Author:** ![system](https://us1.discourse-cdn.com/elastic/original/3X/1/a/1ac57faf039f6b580b3f104ef42a2a89e41014de.png) [@system](https://discuss.elastic.co/u/system)\
**Post date:** [July 6, 2017, 1:27pm UTC](https://discuss.elastic.co/t/aggregation-running-very-slow/28192/9 "2017-07-06T13:27:35Z")

</div>


