# Elasticsearch terms aggregation with partition does not honor the “size” value

**URL:** <https://discuss.elastic.co/t/elasticsearch-terms-aggregation-with-partition-does-not-honor-the-size-value/271294>\
**Category:** Elasticsearch\
**Created:** [April 26, 2021, 7:32pm UTC](https://discuss.elastic.co/t/elasticsearch-terms-aggregation-with-partition-does-not-honor-the-size-value/271294 "2021-04-26T19:32:33Z")\
**Posts on this page:** 6\
**Page:** 1

<div class="post-metadata">

**Author:** ![sayali\_hinge](https://avatars.discourse-cdn.com/v4/letter/s/df788c/32.png) [@sayali\_hinge](https://discuss.elastic.co/u/sayali_hinge)\
**Post date:** [April 26, 2021, 7:32pm UTC](https://discuss.elastic.co/t/elasticsearch-terms-aggregation-with-partition-does-not-honor-the-size-value/271294/1 "2021-04-26T19:32:33Z")

</div>

I have a nested field in the ES index. I used the cardinality aggregation to get the total number of buckets and then came up with num\_partitions, and size values to fetch the buckets.

E.g. the total number of term buckets = 761

Used num\_partitions = 4, size = 200

But this returns me the following number of buckets in the 4 requests = 200, 194, 176, 165 Which sums to 735 \< 761.

Now when using num\_partitions = 3, size = 300

The 3 requests returned the following number of buckets = 246, 261, 254 Which sums to 761.

The second case did not miss any buckets but still, I would expect it to go like this = 300, 300, 161.

So, the questions are -

- Why did the first choice of num\_partitions and size miss 26 buckets?
- Why is not honouring the "size" value in either of the above scenarios?

**NOTE -**

- ES version = 5.6.1. Upgrading to ES 6 is not possible in near future. Hence can not consider using composite aggregation.
- I have referred to [this](https://stackoverflow.com/questions/60411594/elasticsearch-size-value-not-working-in-terms-aggregation-with-partitions) old question and one of the answer quotes the ES documentation that "The terms aggregation is meant to return the top terms and does not allow pagination." But I still do not understand why does it not return buckets in 300, 300, 161 numbers in case 2 above. Without using partitions, it has always honoured the "size" value.

In case required, here is the aggregation part of the query -

```auto
   "aggs": {
       "nested": {
           "path": "software"
       },
       "aggregations": {
           "filtered": {
               "filter": {
                   "bool": {
                       "must": [
                           {
                               "match_phrase_prefix": {
                                   "sw.publisher": {
                                       "query": "O",
                                       "slop": 100,
                                       "max_expansions": 50,
                                       "boost": 1.0
                                   }
                               }
                           }
                       ],
                       "adjust_pure_negative": true,
                       "boost": 1.0
                   }
               },
               "aggregations": {
                   "host_sw": {
                       "terms": {
                           "field": "sw.hostId",
                           "size": 200,
                           "min_doc_count": 1,
                           "shard_min_doc_count": 0,
                           "show_term_doc_count_error": false,
                           "order": [
                               {
                                   "_count": "desc"
                               },
                               {
                                   "_term": "asc"
                               }
                           ],
                           "include": {
                               "partition": 3,
                               "num_partitions": 4
                           }
                       }
                   },
                   "sw_host_id.total_count": {
                       "cardinality": {
                           "field": "sw.hostId",
                           "precision_threshold": 20000
                       }
                   }
               }
           }
       }
   }

```

---

<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 26, 2021, 8:38pm UTC](https://discuss.elastic.co/t/elasticsearch-terms-aggregation-with-partition-does-not-honor-the-size-value/271294/2 "2021-04-26T20:38:43Z")

</div>

> [@sayali\_hinge](#):
>
> Why did the first choice of num\_partitions and size miss 26 buckets?

Because you asked for max 200 results and there were 226 in one of the partitions so not all were returned.

> [@sayali\_hinge](#):
>
> Why is not honouring the "size" value in either of the above scenarios?

The size is the max number of results

What may help explain the behaviour is to understand that partitioning is a rough way to split terms into groups. It is based on hashing the terms and then modulo N so there can be a small degree of unevenness in partition sizes

---

<div class="post-metadata">

**Author:** ![sayali\_hinge](https://avatars.discourse-cdn.com/v4/letter/s/df788c/32.png) [@sayali\_hinge](https://discuss.elastic.co/u/sayali_hinge)\
**Post date:** [April 27, 2021, 8:00am UTC](https://discuss.elastic.co/t/elasticsearch-terms-aggregation-with-partition-does-not-honor-the-size-value/271294/3 "2021-04-27T08:00:37Z")

</div>

THank you @Mark_Harwood , can you please elaborate more on the following?  
Meaning what is N (am I being naive?) and also hashing part.

> [@Mark\_Harwood](#):
>
> It is based on hashing the terms and then modulo N so there can be a small degree of unevenness in partition sizes

Since there can be unevenness, I suppose it can not be used for pagination?  
What I am looking for is a safe formula if I can come up with so that it does not miss any buckets over the pagination.

---

<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 27, 2021, 9:01am UTC](https://discuss.elastic.co/t/elasticsearch-terms-aggregation-with-partition-does-not-honor-the-size-value/271294/4 "2021-04-27T09:01:28Z")

</div>

> [@sayali\_hinge](#):
>
> Meaning what is N (am I being naive?) and also hashing part.

N is the number of partitions you want to break a problem into.  
We use the same algorithm to route a document with an ID to a particular choice of server.  
It's a common technique and [this doc](https://medium.com/system-design-blog/consistent-hashing-b9134c8a9062) discusses it.

The challenge with running aggregations on a distributed data store is you might have many unique values and want to limit the chatter required between data nodes. Take, for example, a store of tweets with many data nodes - each holding a month's worth of tweets.  
Let's say you wanted to query across them all and look at all the Twitter user handles to see who had tweeted about `"covid19"` and how many times they had done that. That is likely to be a lot of unique user accounts so to page through them all you'd need to break the task into multiple requests. If you didn't care about the sort order you could simply page through them alphabetically using the [composite](https://www.elastic.co/guide/en/elasticsearch/reference/current/search-aggregations-bucket-composite-aggregation.html) aggregation with the `after` parameter (rather than the `terms` aggregation). Each of the shards can independently agree that `a` comes before `b` in the alphabet so they can agree without coordination which are the next set of accounts to return for consideration.  
However, if you wanted to look at a list of _the first_ people to tweet about `covid19` the task is much harder because you'd want to sort the account IDs by ascending min date. In a distributed store this is very hard without streaming a lot of data between nodes - unlike an alphabetic sort order the data nodes have no idea if `@DrFauci` or `@JoeShmoe` should be the next candidate to be returned - they don't know for sure who tweeted first about covid 19. What we can do to make the computation simpler is make each request focus on a small subset (aka "partition") of all the millions of unique terms and sort just those by first-tweet-date. If the partitions and size of results returned are sufficiently small the sorted results from each shard for a request can be fused in a way which should be accurate.  
This kind of distributed analytics is made complex in the same way the old [fox, chicken and grain](https://www.mathsisfun.com/chicken_crossing_solution.html) transportation problem is made hard by the constraint of a small boat.

I appreciate there's a lot of complex choices here so I developed a [wizard](https://plnkr.co/edit/eZr7r3KZW02AxNKAHCQa?p=preview) you can run to walk through the options and make a choice.

---

<div class="post-metadata">

**Author:** ![sayali\_hinge](https://avatars.discourse-cdn.com/v4/letter/s/df788c/32.png) [@sayali\_hinge](https://discuss.elastic.co/u/sayali_hinge)\
**Post date:** [April 27, 2021, 11:26am UTC](https://discuss.elastic.co/t/elasticsearch-terms-aggregation-with-partition-does-not-honor-the-size-value/271294/5 "2021-04-27T11:26:55Z")

</div>

Thanks for the detailed response.  
As I mentioned in the question, I can not consider the composite aggregation option yet because we are using ES 5.6 and upgrade is not possible in a near future.

btw, that's a cool tool 'wizard'! but it also is suggesting composite aggregation for my use case.

---

<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 25, 2021, 11:27am UTC](https://discuss.elastic.co/t/elasticsearch-terms-aggregation-with-partition-does-not-honor-the-size-value/271294/6 "2021-05-25T11:27:26Z")

</div>

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