# During terms aggregation with partition we get sum\_other\_doc\_count \> 0 in between partitions

**URL:** <https://discuss.elastic.co/t/during-terms-aggregation-with-partition-we-get-sum-other-doc-count-0-in-between-partitions/321140>\
**Category:** Elasticsearch\
**Created:** [December 13, 2022, 2:40pm UTC](https://discuss.elastic.co/t/during-terms-aggregation-with-partition-we-get-sum-other-doc-count-0-in-between-partitions/321140 "2022-12-13T14:40:53Z")\
**Posts on this page:** 5\
**Page:** 1

<div class="post-metadata">

**Author:** ![akleiber](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/akleiber/32/114637_2.png) [@akleiber](https://discuss.elastic.co/u/akleiber)\
**Post date:** [December 13, 2022, 2:40pm UTC](https://discuss.elastic.co/t/during-terms-aggregation-with-partition-we-get-sum-other-doc-count-0-in-between-partitions/321140/1 "2022-12-13T14:40:53Z")

</div>

We are running ES 6.8 and move some documents to another ES Cluster with a more recent version. We also transform the documents during that process.

We need to group documents by a field, so we do this via a terms aggregation.  
Because there are often more than the search.max\_buckets settings (we can not set it higher) we use the partition option.  
I understand that this is no pagination but for what we want to achieve it is fine that the partitions have different number of results.

> **[Terms Aggregation | Elasticsearch Guide \[6.8\] | Elastic](https://www.elastic.co/guide/en/elasticsearch/reference/6.8/search-aggregations-bucket-terms-aggregation.html#_filtering_values_with_partitions)**

In order to determine how many partitions we need we first run a cardinality aggregation on that field, like suggested in the docs.  
We will then fetch our data with a size of 100, i.e. max 100 buckets for the aggregation per search.

What we see from the results is that `sum_other_doc_count` is sometimes greater than 0. And we wonder what that means in context of using partitioning.

**Our queries**

_Determine number of partitions_

```auto
GET the-index/_search
{
  "aggs": {
    "count": {
      "cardinality": {
        "field": "groupField"
      }
    }
  }, 
  "query": {
    "bool": {
      "filter": {
        "bool": {
          "must": [
            {
              "term": {
                "someField": "someValue"
              }
            }
          ]
        }
      }
    }
  },
  "size": 0
}

```

We now get a number back and divide that bei 100 (our aggregation size setting) so we roughly know how many partitions we need.  
Let's say `cardinality` is 157300 we would set `num_partitions` to `1573`.

_Fetch buckets_

```auto
GET the-index/_search
{
  "aggs": {
    "byGroupField": {
      "terms": {
        "execution_hint": "map",
        "field": "groupField",
        "size": 100,
        "include": {
          "partition": 123,
          "num_partitions": 1573
        }
      }
    }
  },
  "query": {
    "bool": {
      "filter": {
        "bool": {
          "must": [
            {
              "term": {
                "someField": "someValue"
              }
            }
          ]
        }
      }
    }
  },
  "size": 0
}

```

What we observed now is `sum_other_doc_count` being greater than 0 for some partitions. Even when we set `num_partitions` way higher than necessary (in respect to the size setting) we will get `sum_other_doc_count` greater 0 for many partitions.

What we expected was that `sum_other_doc_count` is greater 0 only for the last partition in case our `size` for the aggregation is too low.

Can somebody point me to some docs or post an explanation why we get `sum_other_doc_count` greater 0 for many partitions?

Thank you very much 🙂

---

<div class="post-metadata">

**Author:** ![Mark\_Harwood1](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mark_harwood1/32/101255_2.png) [@Mark\_Harwood1](https://discuss.elastic.co/u/Mark_Harwood1)\
**Post date:** [December 13, 2022, 3:04pm UTC](https://discuss.elastic.co/t/during-terms-aggregation-with-partition-we-get-sum-other-doc-count-0-in-between-partitions/321140/2 "2022-12-13T15:04:51Z")

</div>

> [@akleiber](#):
>
> Even when we set `num_partitions` way higher than necessary (in respect to the size setting) we will get `sum_other_doc_count` greater 0 for many partitions.

What were you picking for num\_partitions?

> [@akleiber](#):
>
> Can somebody point me to some docs or post an explanation why we get `sum_other_doc_count` greater 0 for many partitions?

The routing for which terms go in which partitions is based on hash-modulo eg. term.hashcode()%num\_partitions  
Being based on a hashing algorithm we can expect some unevenness in numbers of terms in each partition (e.g. it might be 95, 99, 96, 100, 94, 102......)  
That's why when you set size to 100 the last partition in the example above would have 2 missing terms.

---

<div class="post-metadata">

**Author:** ![akleiber](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/akleiber/32/114637_2.png) [@akleiber](https://discuss.elastic.co/u/akleiber)\
**Post date:** [December 13, 2022, 3:09pm UTC](https://discuss.elastic.co/t/during-terms-aggregation-with-partition-we-get-sum-other-doc-count-0-in-between-partitions/321140/3 "2022-12-13T15:09:37Z")

</div>

> [@Mark\_Harwood1](#):
>
> What were you picking for num\_partitions?

1773

> [@Mark\_Harwood1](#):
>
> That's why when you set size to 100 the last partition in the example above would have 2 missing terms.

Ah thank you for that explanation. And those terms will also not be included in the next partition right? So we need to choose a "better" size in respect to `num_partitions` ?  
I am still a bit confused how to avoid tose missing? terms.

---

<div class="post-metadata">

**Author:** ![Mark\_Harwood1](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mark_harwood1/32/101255_2.png) [@Mark\_Harwood1](https://discuss.elastic.co/u/Mark_Harwood1)\
**Post date:** [December 13, 2022, 3:29pm UTC](https://discuss.elastic.co/t/during-terms-aggregation-with-partition-we-get-sum-other-doc-count-0-in-between-partitions/321140/4 "2022-12-13T15:29:10Z")

</div>

> [@akleiber](#):
>
> Ah thank you for that explanation. And those terms will also not be included in the next partition right?

Correct. Hash-modulo routing is a simple and efficient way to organise potentially billions of values into (roughly) equal sized groups. Given the groups can vary a little above and below your target size just set the retrieval 'size' setting to something \>100 to allow for this overspill. Target-partition-size x 2 (e.g. 100 x 2 = 200) should be more than enough I'd have thought to compensate for hashing variations in partition size.

I notice your example is only sorting the terms by their value and not by anything more complex like a derived sum of sales. In the simpler cases the [composite aggregation](https://www.elastic.co/guide/en/elasticsearch/reference/current/search-aggregations-bucket-composite-aggregation.html#_pagination) would be a better way to page through results

---

<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:** [January 10, 2023, 3:29pm UTC](https://discuss.elastic.co/t/during-terms-aggregation-with-partition-we-get-sum-other-doc-count-0-in-between-partitions/321140/5 "2023-01-10T15:29:30Z")

</div>

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