Spec-Zone.ru › Elasticsearch 8
›Elasticsearch Guide [8.17] ›Агрегирования

Агрегирования типа Pipeline

Агрегирования типа Pipeline работают с результатами других агрегирований, а не с наборами документов, добавляя информацию в дерево результатов. Существует множество различных типов агрегирований Pipeline, каждый из которых вычисляет разную информацию из других агрегирований, но эти типы можно разделить на две группы:

Родительский
Семейство агрегирований Pipeline, которое получает результаты родительского агрегирования и способно вычислять новые корзины или новые агрегирования для добавления в существующие корзины.
Братский
Агрегирования Pipeline, которые получают результаты братского агрегирования и способны вычислить новое агрегирование, которое будет на том же уровне, что и братское агрегирование.

Агрегирования Pipeline могут ссылаться на агрегирования, необходимые для вычисления, используя параметр buckets_path для указания путей к требуемым метрикам. Синтаксис определения этих путей можно найти в разделе buckets_path Синтаксис ниже.

Агрегирования Pipeline не могут иметь вложенных агрегирований, но в зависимости от типа они могут ссылаться на другое агрегирование в buckets_path, позволяя цеплять агрегирования Pipeline. Например, вы можете объединить два производных, чтобы вычислить вторую производную (т.е. производную от производной).

Поскольку агрегирования Pipeline только добавляют к выводу, при объединении агрегирований Pipeline выход каждого агрегирования Pipeline будет включен в конечный вывод.

buckets_path Синтаксис

Большинство агрегирований Pipeline требуют другого агрегирования в качестве входных данных. Входное агрегирование определяется с помощью параметра buckets_path, который следует определенному формату:

AGG_SEPARATOR       =  `>` ;
METRIC_SEPARATOR    =  `.` ;
AGG_NAME            =  <the name of the aggregation> ;
METRIC              =  <the name of the metric (in case of multi-value metrics aggregation)> ;
MULTIBUCKET_KEY     =  `[<KEY_NAME>]`
PATH                =  <AGG_NAME><MULTIBUCKET_KEY>? (<AGG_SEPARATOR>, <AGG_NAME> )* ( <METRIC_SEPARATOR>, <METRIC> ) ;

Например, путь "my_bucket>my_stats.avg" приведет к значению avg в метрике "my_stats", которая содержится в агрегировании корзины "my_bucket".

Вот еще несколько примеров:

  • multi_bucket["foo"]>single_bucket>multi_metric.avg перейдет к метрике avg в агрегировании "multi_metric" под единственной корзинной "single_bucket" в корзинном агрегировании "foo" многокорзинного агрегирования "multi_bucket".
  • agg1["foo"]._count получит метрику _count для корзины "foo" в многокорзинном агрегировании "multi_bucket"

Пути относительны к положению агрегирования Pipeline; они не являются абсолютными путями, и путь не может возвращаться "вверх" по дереву агрегирования. Например, эта производная встроена в date_histogram и ссылается на "братскую" метрику "the_sum":

resp = client.search(
    aggs={
        "my_date_histo": {
            "date_histogram": {
                "field": "timestamp",
                "calendar_interval": "day"
            },
            "aggs": {
                "the_sum": {
                    "sum": {
                        "field": "lemmings"
                    }
                },
                "the_deriv": {
                    "derivative": {
                        "buckets_path": "the_sum"
                    }
                }
            }
        }
    },
)
print(resp)
response = client.search(
  body: {
    aggregations: {
      my_date_histo: {
        date_histogram: {
          field: 'timestamp',
          calendar_interval: 'day'
        },
        aggregations: {
          the_sum: {
            sum: {
              field: 'lemmings'
            }
          },
          the_deriv: {
            derivative: {
              buckets_path: 'the_sum'
            }
          }
        }
      }
    }
  }
)
puts response
const response = await client.search({
  aggs: {
    my_date_histo: {
      date_histogram: {
        field: "timestamp",
        calendar_interval: "day",
      },
      aggs: {
        the_sum: {
          sum: {
            field: "lemmings",
          },
        },
        the_deriv: {
          derivative: {
            buckets_path: "the_sum",
          },
        },
      },
    },
  },
});
console.log(response);
POST /_search
{
  "aggs": {
    "my_date_histo": {
      "date_histogram": {
        "field": "timestamp",
        "calendar_interval": "day"
      },
      "aggs": {
        "the_sum": {
          "sum": { "field": "lemmings" }              
        },
        "the_deriv": {
          "derivative": { "buckets_path": "the_sum" } 
        }
      }
    }
  }
}

Метрика называется "the_sum"

buckets_path ссылается на метрику по относительному пути "the_sum"

buckets_path также используется для братских агрегирований Pipeline, где агрегирование находится "рядом" с рядом корзин, а не "внутри" их. Например, агрегирование max_bucket использует buckets_path для указания метрики, встроенной внутри братского агрегирования:

resp = client.search(
    aggs={
        "sales_per_month": {
            "date_histogram": {
                "field": "date",
                "calendar_interval": "month"
            },
            "aggs": {
                "sales": {
                    "sum": {
                        "field": "price"
                    }
                }
            }
        },
        "max_monthly_sales": {
            "max_bucket": {
                "buckets_path": "sales_per_month>sales"
            }
        }
    },
)
print(resp)
response = client.search(
  body: {
    aggregations: {
      sales_per_month: {
        date_histogram: {
          field: 'date',
          calendar_interval: 'month'
        },
        aggregations: {
          sales: {
            sum: {
              field: 'price'
            }
          }
        }
      },
      max_monthly_sales: {
        max_bucket: {
          buckets_path: 'sales_per_month>sales'
        }
      }
    }
  }
)
puts response
const response = await client.search({
  aggs: {
    sales_per_month: {
      date_histogram: {
        field: "date",
        calendar_interval: "month",
      },
      aggs: {
        sales: {
          sum: {
            field: "price",
          },
        },
      },
    },
    max_monthly_sales: {
      max_bucket: {
        buckets_path: "sales_per_month>sales",
      },
    },
  },
});
console.log(response);
POST /_search
{
  "aggs": {
    "sales_per_month": {
      "date_histogram": {
        "field": "date",
        "calendar_interval": "month"
      },
      "aggs": {
        "sales": {
          "sum": {
            "field": "price"
          }
        }
      }
    },
    "max_monthly_sales": {
      "max_bucket": {
        "buckets_path": "sales_per_month>sales" 
      }
    }
  }
}

buckets_path сообщает агрегированию max_bucket, что мы хотим максимальное значение агрегирования sales в date histogram sales_per_month.

Если братское агрегирование Pipeline ссылается на многокорзинное агрегирование, такое как terms agg, у него также есть возможность выбрать определенные ключи из многокорзинного агрегирования. Например, bucket_script может выбрать две определенные корзины (по их ключам корзины) для выполнения вычисления:

resp = client.search(
    aggs={
        "sales_per_month": {
            "date_histogram": {
                "field": "date",
                "calendar_interval": "month"
            },
            "aggs": {
                "sale_type": {
                    "terms": {
                        "field": "type"
                    },
                    "aggs": {
                        "sales": {
                            "sum": {
                                "field": "price"
                            }
                        }
                    }
                },
                "hat_vs_bag_ratio": {
                    "bucket_script": {
                        "buckets_path": {
                            "hats": "sale_type['hat']>sales",
                            "bags": "sale_type['bag']>sales"
                        },
                        "script": "params.hats / params.bags"
                    }
                }
            }
        }
    },
)
print(resp)
response = client.search(
  body: {
    aggregations: {
      sales_per_month: {
        date_histogram: {
          field: 'date',
          calendar_interval: 'month'
        },
        aggregations: {
          sale_type: {
            terms: {
              field: 'type'
            },
            aggregations: {
              sales: {
                sum: {
                  field: 'price'
                }
              }
            }
          },
          hat_vs_bag_ratio: {
            bucket_script: {
              buckets_path: {
                hats: "sale_type['hat']>sales",
                bags: "sale_type['bag']>sales"
              },
              script: 'params.hats / params.bags'
            }
          }
        }
      }
    }
  }
)
puts response
const response = await client.search({
  aggs: {
    sales_per_month: {
      date_histogram: {
        field: "date",
        calendar_interval: "month",
      },
      aggs: {
        sale_type: {
          terms: {
            field: "type",
          },
          aggs: {
            sales: {
              sum: {
                field: "price",
              },
            },
          },
        },
        hat_vs_bag_ratio: {
          bucket_script: {
            buckets_path: {
              hats: "sale_type['hat']>sales",
              bags: "sale_type['bag']>sales",
            },
            script: "params.hats / params.bags",
          },
        },
      },
    },
  },
});
console.log(response);
POST /_search
{
  "aggs": {
    "sales_per_month": {
      "date_histogram": {
        "field": "date",
        "calendar_interval": "month"
      },
      "aggs": {
        "sale_type": {
          "terms": {
            "field": "type"
          },
          "aggs": {
            "sales": {
              "sum": {
                "field": "price"
              }
            }
          }
        },
        "hat_vs_bag_ratio": {
          "bucket_script": {
            "buckets_path": {
              "hats": "sale_type['hat']>sales",   
              "bags": "sale_type['bag']>sales"    
            },
            "script": "params.hats / params.bags"
          }
        }
      }
    }
  }
}

buckets_path выбирает корзины hats и bags (через ['hat']/['bag']`) для использования в скрипте, а не извлекает все корзины из агрегирования sale_type

Специальные пути

Вместо указания пути к метрике, buckets_path может использовать специальный путь "_count". Это указывает агрегированию Pipeline использовать количество документов в качестве входных данных. Например, производная может быть вычислена по количеству документов в каждой корзине, а не по конкретной метрике:

resp = client.search(
    aggs={
        "my_date_histo": {
            "date_histogram": {
                "field": "timestamp",
                "calendar_interval": "day"
            },
            "aggs": {
                "the_deriv": {
                    "derivative": {
                        "buckets_path": "_count"
                    }
                }
            }
        }
    },
)
print(resp)
response = client.search(
  body: {
    aggregations: {
      my_date_histo: {
        date_histogram: {
          field: 'timestamp',
          calendar_interval: 'day'
        },
        aggregations: {
          the_deriv: {
            derivative: {
              buckets_path: '_count'
            }
          }
        }
      }
    }
  }
)
puts response
const response = await client.search({
  aggs: {
    my_date_histo: {
      date_histogram: {
        field: "timestamp",
        calendar_interval: "day",
      },
      aggs: {
        the_deriv: {
          derivative: {
            buckets_path: "_count",
          },
        },
      },
    },
  },
});
console.log(response);
POST /_search
{
  "aggs": {
    "my_date_histo": {
      "date_histogram": {
        "field": "timestamp",
        "calendar_interval": "day"
      },
      "aggs": {
        "the_deriv": {
          "derivative": { "buckets_path": "_count" } 
        }
      }
    }
  }
}

Используя _count вместо имени метрики, мы можем вычислить производную от количества документов в гистограмме

buckets_path также может использовать "_bucket_count" и путь к многокорзинному агрегированию, чтобы использовать количество корзин, возвращенных этим агрегированием, в агрегировании Pipeline, а не метрику. Например, здесь может быть использован bucket_selector для фильтрации корзин, которые не содержат корзин для внутреннего агрегирования terms:

resp = client.search(
    index="sales",
    size=0,
    aggs={
        "histo": {
            "date_histogram": {
                "field": "date",
                "calendar_interval": "day"
            },
            "aggs": {
                "categories": {
                    "terms": {
                        "field": "category"
                    }
                },
                "min_bucket_selector": {
                    "bucket_selector": {
                        "buckets_path": {
                            "count": "categories._bucket_count"
                        },
                        "script": {
                            "source": "params.count != 0"
                        }
                    }
                }
            }
        }
    },
)
print(resp)
response = client.search(
  index: 'sales',
  body: {
    size: 0,
    aggregations: {
      histo: {
        date_histogram: {
          field: 'date',
          calendar_interval: 'day'
        },
        aggregations: {
          categories: {
            terms: {
              field: 'category'
            }
          },
          min_bucket_selector: {
            bucket_selector: {
              buckets_path: {
                count: 'categories._bucket_count'
              },
              script: {
                source: 'params.count != 0'
              }
            }
          }
        }
      }
    }
  }
)
puts response
const response = await client.search({
  index: "sales",
  size: 0,
  aggs: {
    histo: {
      date_histogram: {
        field: "date",
        calendar_interval: "day",
      },
      aggs: {
        categories: {
          terms: {
            field: "category",
          },
        },
        min_bucket_selector: {
          bucket_selector: {
            buckets_path: {
              count: "categories._bucket_count",
            },
            script: {
              source: "params.count != 0",
            },
          },
        },
      },
    },
  },
});
console.log(response);
POST /sales/_search
{
  "size": 0,
  "aggs": {
    "histo": {
      "date_histogram": {
        "field": "date",
        "calendar_interval": "day"
      },
      "aggs": {
        "categories": {
          "terms": {
            "field": "category"
          }
        },
        "min_bucket_selector": {
          "bucket_selector": {
            "buckets_path": {
              "count": "categories._bucket_count" 
            },
            "script": {
              "source": "params.count != 0"
            }
          }
        }
      }
    }
  }
}

Используя _bucket_count вместо имени метрики, мы можем отфильтровать корзины histo, где они не содержат корзин для агрегирования categories

Обработка точек в именах агрегирований

Поддерживается альтернативный синтаксис для работы с агрегированиями или метриками, имеющими точки в имени, например, 99.9-й перцентиль. Эта метрика может быть указана как:

"buckets_path": "my_percentile[99.9]"

Обработка пробелов в данных

Данные в реальном мире часто шумные и иногда содержат пробелы — места, где данных просто нет. Это может происходить по разным причинам, наиболее распространенными из которых являются:

  • Документы, попадающие в корзину, не содержат требуемого поля
  • Нет документов, соответствующих запросу для одной или нескольких корзин
  • Метрика, которая вычисляется, не может сгенерировать значение, вероятно, потому что у другой зависимой корзины отсутствует значение. Некоторые агрегирования Pipeline имеют определенные требования, которые должны быть выполнены (например, производная не может вычислить метрику для первого значения, потому что нет предыдущего значения, скользящее среднее HoltWinters требует начальных данных для начала вычисления и т.д.)

Политики обработки пробелов — это механизм, позволяющий агрегированию Pipeline узнать о желаемом поведении при столкновении с "пробелами" или отсутствием данных. Все агрегирования Pipeline принимают параметр gap_policy. В настоящее время доступны две политики обработки пробелов:

пропустить
Этот вариант рассматривает отсутствующие данные так, как будто корзины не существует. Он пропустит корзину и продолжит вычисления, используя следующее доступное значение.
вставить_ноль
Этот вариант заменит отсутствующие значения нулем (0), и вычисление агрегирования Pipeline продолжится как обычно.
сохранить_значения
Этот вариант похож на пропуск, за исключением того, что если метрика предоставляет не нулевое, не NaN значение, это значение используется, в противном случае пустая корзина пропускается.

© 2023-2025 Elasticsearch
As of September 2024, Elasticsearch is available under a choice of three licenses: the Server Side Public License (SSPL), the Elastic License, or the AGPLv3 (OSI approved).
Elasticsearch and the Elasticsearch logo are trademarks of Elasticsearch B.V., registered in the U.S. and in other countries.
https://www.elastic.co/guide/en/elasticsearch/reference/8.17/search-aggregations-pipeline.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API