# Sum Aggregation returning bad value

**URL:** <https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538>\
**Category:** Elasticsearch\
**Created:** [April 5, 2019, 8:06am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538 "2019-04-05T08:06:30Z")\
**Posts on this page:** 11\
**Page:** 1

<div class="post-metadata">

**Author:** ![Nerijus\_Oftas](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/nerijus_oftas/32/32616_2.png) [@Nerijus\_Oftas](https://discuss.elastic.co/u/Nerijus_Oftas)\
**Post date:** [April 5, 2019, 8:06am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/1 "2019-04-05T08:06:31Z")

</div>

Hi,  
I am using a simple sum aggregation query. It's fine when I have a small count of documents. But if the query needs to aggregate from 10 million filtered documents, the sum value is not correct. I found that the solution is to increase shard\_size parameter. But this is a temporary solution. Maybe there is a better solution to do SUM aggregation?  
Elastic search version 5.5  
My query:

```
       {
                  "query": {
                    "bool": {
                      "must": [
                        {
                          "terms": {
                            "round_id": ["1","2","3"]
                          }
                        }
                      ],
                      "must_not": [
                        {
                          "term": {
                            "user_id": 0
                        }
                      }
                      ]
                    }
                  },
                  "size": 0,
                  "aggs": {
                    "groups": {
                      "terms": {
                        "field": "user_id",
                        "size": 10,
                        "order": {"total_points_sum": "desc"}
                      },
                      "aggs": {
                        "total_points_sum": {
                          "sum": {
                            "field": "points"
                          }
                        }
                      }
                    }
                  }
                }

```

This query gets more than 10 million documents. I need only 10 aggregated results with the biggest sum value.  
I would be very grateful if someone could help me.

---

<div class="post-metadata">

**Author:** ![Mark\_Harwood](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mark_harwood/32/10538_2.png) [@Mark\_Harwood](https://discuss.elastic.co/u/Mark_Harwood)\
**Post date:** [April 5, 2019, 8:38am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/2 "2019-04-05T08:38:41Z")

</div>

Hi Nerijus,

> [@Nerijus\_Oftas](#):
>
> . I found that the solution is to increase shard\_size parameter. But this is a temporary solution.

Why do you say increasing shard\_size is only a temporary solution?  
We can talk about alternative strategies like term partitioning or entity-centric indexes but first it would be good to know if it's important that the top 10 users are the true top 10 (all things considered) or that the sums shown for the (maybe slightly inaccurate) top 10 are accurate numbers for those users. If the latter then you just need to run a follow-up request where your `terms` agg lists the user-ids in the `includes` clause.

---

<div class="post-metadata">

**Author:** ![Nerijus\_Oftas](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/nerijus_oftas/32/32616_2.png) [@Nerijus\_Oftas](https://discuss.elastic.co/u/Nerijus_Oftas)\
**Post date:** [April 5, 2019, 9:09am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/3 "2019-04-05T09:09:57Z")

</div>

I think that shard\_size is temp solution because if I will assign the value lets to say 10000 after some time when I will have more documents again I will need to increase shard\_size. I am not sure if it's ok to add very big shard\_size value. How much we can increase it? At this moment we see that increased shard\_size helps to get not approximately data.  
The 10 users I get are not the same which I should get with the correct sum values. So I can not run follow-up with these 10 users ids.  
Maybe it's possible to write the query in another way?

---

<div class="post-metadata">

**Author:** ![Mark\_Harwood](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mark_harwood/32/10538_2.png) [@Mark\_Harwood](https://discuss.elastic.co/u/Mark_Harwood)\
**Post date:** [April 5, 2019, 9:27am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/4 "2019-04-05T09:27:21Z")

</div>

> [@Nerijus\_Oftas](#):
>
> I am not sure if it's ok to add very big shard\_size value.

Given you're asking for the biggest users (as opposed to the smallest sums) the shard\_size required should be manageable, meaning it should fit into RAM used for a single request).  
The way to think about this is that querying many shards for the top N of something is like asking a group of people for their favourite album in order to find their shared favourite. This question is effectively shard\_size=1 and could lead to an inaccurate result (every person in the group returns a personal favourite that is not shared by any other member of the group). If you asked each person for the top 10 albums (shard\_size=10) then you might discover the group's favourite album was Nirvana's NeverMind, (it being ranked 5th, 7th and 9th by 3 different people).  
In theory, each person might have such differing tastes that you'd have to ask them for their top _million_ albums before you'd discover anything they have in common. That's _theoretically_ possible but in practice highly unlikely. The same is true if you randomly distribute data across shards.

We report back on the bounds of the inaccuracies in results so you can adjust shard\_size accordingly. If you do hit memory errors with large shard\_sizes then you can either:  
A) Break your query into multiple requests using the partitioning feature for `terms` aggs or  
B) Bring related data closer to each other at index time using routing (ensures same shard) or entity-centric indexing (same document).

---

<div class="post-metadata">

**Author:** ![Nerijus\_Oftas](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/nerijus_oftas/32/32616_2.png) [@Nerijus\_Oftas](https://discuss.elastic.co/u/Nerijus_Oftas)\
**Post date:** [April 6, 2019, 7:16am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/5 "2019-04-06T07:16:46Z")

</div>

Thanks for your answers. To be sure if I understood correctly.  
A) is partitioning feature actually same like scroll but its for aggregation? so we need to run lot of queries and just increase partition number? How I understand this works we need to count unique users id using cardinality aggregation for num\_partitions. But cardinality is also approximately number. how to get exact number to know partition number count then? [https://www.elastic.co/guide/en/elasticsearch/reference/master/search-aggregations-bucket-terms-aggregation.html#\_filtering\_values\_with\_partitions](https://www.elastic.co/guide/en/elasticsearch/reference/master/search-aggregations-bucket-terms-aggregation.html#_filtering_values_with_partitions)  
B) its had to find more info about entity-centric indexing ... And I did not get how it works. Maybe you can provide link to documentation about this?

---

<div class="post-metadata">

**Author:** ![Mark\_Harwood](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mark_harwood/32/10538_2.png) [@Mark\_Harwood](https://discuss.elastic.co/u/Mark_Harwood)\
**Post date:** [April 6, 2019, 8:38am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/6 "2019-04-06T08:38:09Z")

</div>

> [@Nerijus\_Oftas](#):
>
> is partitioning feature actually same like scroll

It ensures that all shards look at the same arbitrary subset of terms. For each term we compute hashcode modulo N and if the answer comes out as the chosen partition number then independent shards can agree that this term is one of the terms we are considering in this pass over the data. By reducing the number of unique terms considered in any one request we can see the reported error margins decrease towards zero as we are able to pack fully complete details for considered terms into the space-limited shard responses

For entity centric indexing see talk and example scripts/data here: [https://twitter.com/elasticmark/status/1009380268409610240?s=21](https://twitter.com/elasticmark/status/1009380268409610240?s=21)

---

<div class="post-metadata">

**Author:** ![Nerijus\_Oftas](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/nerijus_oftas/32/32616_2.png) [@Nerijus\_Oftas](https://discuss.elastic.co/u/Nerijus_Oftas)\
**Post date:** [April 8, 2019, 6:26am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/7 "2019-04-08T06:26:22Z")

</div>

Thanks for your replay.  
So to resume the best way how I understand get accurate results not approximate is to increasing "shard\_size".  
And the last things:

1. From your experience what is the best way to calculate "shard\_size" size. For example instance is 32 Memory (GiB).
2. Now for example we have about 25 millions records in one index and use "SUM" query like wrote on first post. If we will split to the smaller index (daily index) will it be a big impact to get the better accuracy of the results?  
Thanks.

---

<div class="post-metadata">

**Author:** ![Mark\_Harwood](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mark_harwood/32/10538_2.png) [@Mark\_Harwood](https://discuss.elastic.co/u/Mark_Harwood)\
**Post date:** [April 8, 2019, 5:05pm UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/8 "2019-04-08T17:05:50Z")

</div>

> [@Nerijus\_Oftas](#):
>
> what is the best way to calculate "shard\_size" size.

The lowest number that produces zero reported error margin in the results - plus a factor to allow for future growth.

> [@Nerijus\_Oftas](#):
>
> will it be a big impact to get the better accuracy of the results?

Creating more division between related content will only exacerbate the accuracy concerns. It is precisely because you don’t keep all data in the same place that we have accuracy issues. We use distributed indices to deal with scale or ease deletion of old data, not improve accuracy.

---

<div class="post-metadata">

**Author:** ![Nerijus\_Oftas](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/nerijus_oftas/32/32616_2.png) [@Nerijus\_Oftas](https://discuss.elastic.co/u/Nerijus_Oftas)\
**Post date:** [April 9, 2019, 5:48am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/9 "2019-04-09T05:48:11Z")

</div>

Thank you for your patience and answers.

---

<div class="post-metadata">

**Author:** ![Mark\_Harwood](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mark_harwood/32/10538_2.png) [@Mark\_Harwood](https://discuss.elastic.co/u/Mark_Harwood)\
**Post date:** [April 9, 2019, 6:27am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/10 "2019-04-09T06:27:34Z")

</div>

You’re very welcome

---

<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:** [May 7, 2019, 6:27am UTC](https://discuss.elastic.co/t/sum-aggregation-returning-bad-value/175538/11 "2019-05-07T06:27:38Z")

</div>

This topic was automatically closed 28 days after the last reply. New replies are no longer allowed.
