Spec-Zone.ru › Elasticsearch 8
›Elasticsearch Guide [8.17] ›Свертка или преобразование данных ›Преобразование данных

Примеры преобразования

Эти примеры демонстрируют, как использовать преобразования, чтобы извлечь полезные сведения из ваших данных. Все примеры используют один из наборов образцовых данных Kibana. Для более подробного пошагового примера см. Учебник: Преобразование образцовых данных электронной коммерции.

  • Нахождение лучших клиентов
  • Нахождение авиаперевозчиков с наибольшим количеством задержек
  • Нахождение подозрительных IP-адресов клиентов
  • Нахождение последнего события регистрации для каждого IP-адреса
  • Нахождение IP-адресов клиентов, которые отправили наибольшее количество байтов на сервер
  • Получение имени и адреса электронной почты клиента по идентификатору клиента

Нахождение лучших клиентов

В этом примере используется набор образцовых данных заказов электронной коммерции для поиска клиентов, потративших больше всего в гипотетическом интернет-магазине. Давайте воспользуемся преобразованием типа pivot, таким образом, чтобы целевой индекс содержал количество заказов, общую стоимость заказов, количество уникальных продуктов, среднюю цену за заказ и общее количество заказанных продуктов для каждого клиента.

Finding your best customers with transforms in Kibana

В качестве альтернативы можно использовать преобразование предварительного просмотра и API-интерфейс создания преобразований.

Пример API
resp = client.transform.preview_transform(
    source={
        "index": "kibana_sample_data_ecommerce"
    },
    dest={
        "index": "sample_ecommerce_orders_by_customer"
    },
    pivot={
        "group_by": {
            "user": {
                "terms": {
                    "field": "user"
                }
            },
            "customer_id": {
                "terms": {
                    "field": "customer_id"
                }
            }
        },
        "aggregations": {
            "order_count": {
                "value_count": {
                    "field": "order_id"
                }
            },
            "total_order_amt": {
                "sum": {
                    "field": "taxful_total_price"
                }
            },
            "avg_amt_per_order": {
                "avg": {
                    "field": "taxful_total_price"
                }
            },
            "avg_unique_products_per_order": {
                "avg": {
                    "field": "total_unique_products"
                }
            },
            "total_unique_products": {
                "cardinality": {
                    "field": "products.product_id"
                }
            }
        }
    },
)
print(resp)
const response = await client.transform.previewTransform({
  source: {
    index: "kibana_sample_data_ecommerce",
  },
  dest: {
    index: "sample_ecommerce_orders_by_customer",
  },
  pivot: {
    group_by: {
      user: {
        terms: {
          field: "user",
        },
      },
      customer_id: {
        terms: {
          field: "customer_id",
        },
      },
    },
    aggregations: {
      order_count: {
        value_count: {
          field: "order_id",
        },
      },
      total_order_amt: {
        sum: {
          field: "taxful_total_price",
        },
      },
      avg_amt_per_order: {
        avg: {
          field: "taxful_total_price",
        },
      },
      avg_unique_products_per_order: {
        avg: {
          field: "total_unique_products",
        },
      },
      total_unique_products: {
        cardinality: {
          field: "products.product_id",
        },
      },
    },
  },
});
console.log(response);
POST _transform/_preview
{
  "source": {
    "index": "kibana_sample_data_ecommerce"
  },
  "dest" : { 
    "index" : "sample_ecommerce_orders_by_customer"
  },
  "pivot": {
    "group_by": { 
      "user": { "terms": { "field": "user" }},
      "customer_id": { "terms": { "field": "customer_id" }}
    },
    "aggregations": {
      "order_count": { "value_count": { "field": "order_id" }},
      "total_order_amt": { "sum": { "field": "taxful_total_price" }},
      "avg_amt_per_order": { "avg": { "field": "taxful_total_price" }},
      "avg_unique_products_per_order": { "avg": { "field": "total_unique_products" }},
      "total_unique_products": { "cardinality": { "field": "products.product_id" }}
    }
  }
}

Целевой индекс для преобразования. Он игнорируется _preview.

Выбраны два group_by поля. Это означает, что преобразование содержит уникальную строку на сочетание user и customer_id. В этом наборе данных оба поля уникальны. Включение обоих полей в преобразование обеспечивает больший контекст для конечных результатов.

В примере выше используется сжатый формат JSON для лучшей читабельности объекта pivot.

API предварительного просмотра преобразований позволяет предварительно просмотреть макет преобразования, заполненный некоторыми образцовыми значениями. Например:

{
  "preview" : [
    {
      "total_order_amt" : 3946.9765625,
      "order_count" : 59.0,
      "total_unique_products" : 116.0,
      "avg_unique_products_per_order" : 2.0,
      "customer_id" : "10",
      "user" : "recip",
      "avg_amt_per_order" : 66.89790783898304
    },
    ...
    ]
  }

Это преобразование упрощает ответы на такие вопросы, как:

  • Какие клиенты тратят больше всего?
  • Какие клиенты тратят больше всего на заказ?
  • Какие клиенты чаще всего совершают заказы?
  • Какие клиенты заказали наименьшее количество различных продуктов?

Эти вопросы можно ответить с помощью только агрегаций, однако преобразования позволяют нам сохранить эти данные в индексе, ориентированном на клиента. Это позволяет нам анализировать данные в масштабе и обеспечивает большую гибкость при исследовании и навигации по данным с точки зрения клиента. В некоторых случаях это может даже упростить создание визуализаций.

Нахождение авиаперевозчиков с наибольшим количеством задержек

В этом примере используется набор образцовых данных рейсов, чтобы определить, какая авиакомпания имела больше всего задержек. Сначала отфильтруйте исходные данные таким образом, чтобы исключить все отменённые рейсы с помощью фильтра запроса. Затем преобразуйте данные так, чтобы они содержали количество уникальных рейсов, сумму задержанных минут и сумму минут полета по авиакомпании. Наконец, используйте bucket_script, чтобы определить, какой процент времени полета фактически был задержан.

resp = client.transform.preview_transform(
    source={
        "index": "kibana_sample_data_flights",
        "query": {
            "bool": {
                "filter": [
                    {
                        "term": {
                            "Cancelled": False
                        }
                    }
                ]
            }
        }
    },
    dest={
        "index": "sample_flight_delays_by_carrier"
    },
    pivot={
        "group_by": {
            "carrier": {
                "terms": {
                    "field": "Carrier"
                }
            }
        },
        "aggregations": {
            "flights_count": {
                "value_count": {
                    "field": "FlightNum"
                }
            },
            "delay_mins_total": {
                "sum": {
                    "field": "FlightDelayMin"
                }
            },
            "flight_mins_total": {
                "sum": {
                    "field": "FlightTimeMin"
                }
            },
            "delay_time_percentage": {
                "bucket_script": {
                    "buckets_path": {
                        "delay_time": "delay_mins_total.value",
                        "flight_time": "flight_mins_total.value"
                    },
                    "script": "(params.delay_time / params.flight_time) * 100"
                }
            }
        }
    },
)
print(resp)
const response = await client.transform.previewTransform({
  source: {
    index: "kibana_sample_data_flights",
    query: {
      bool: {
        filter: [
          {
            term: {
              Cancelled: false,
            },
          },
        ],
      },
    },
  },
  dest: {
    index: "sample_flight_delays_by_carrier",
  },
  pivot: {
    group_by: {
      carrier: {
        terms: {
          field: "Carrier",
        },
      },
    },
    aggregations: {
      flights_count: {
        value_count: {
          field: "FlightNum",
        },
      },
      delay_mins_total: {
        sum: {
          field: "FlightDelayMin",
        },
      },
      flight_mins_total: {
        sum: {
          field: "FlightTimeMin",
        },
      },
      delay_time_percentage: {
        bucket_script: {
          buckets_path: {
            delay_time: "delay_mins_total.value",
            flight_time: "flight_mins_total.value",
          },
          script: "(params.delay_time / params.flight_time) * 100",
        },
      },
    },
  },
});
console.log(response);
POST _transform/_preview
{
  "source": {
    "index": "kibana_sample_data_flights",
    "query": { 
      "bool": {
        "filter": [
          { "term":  { "Cancelled": false } }
        ]
      }
    }
  },
  "dest" : { 
    "index" : "sample_flight_delays_by_carrier"
  },
  "pivot": {
    "group_by": { 
      "carrier": { "terms": { "field": "Carrier" }}
    },
    "aggregations": {
      "flights_count": { "value_count": { "field": "FlightNum" }},
      "delay_mins_total": { "sum": { "field": "FlightDelayMin" }},
      "flight_mins_total": { "sum": { "field": "FlightTimeMin" }},
      "delay_time_percentage": { 
        "bucket_script": {
          "buckets_path": {
            "delay_time": "delay_mins_total.value",
            "flight_time": "flight_mins_total.value"
          },
          "script": "(params.delay_time / params.flight_time) * 100"
        }
      }
    }
  }
}

Отфильтруйте исходные данные, чтобы выбрать только рейсы, которые не отменены.

Целевой индекс для преобразования. Он игнорируется _preview.

Данные сгруппированы по полю Carrier, которое содержит название авиакомпании.

Это bucket_script выполняет вычисления над результатами, возвращаемыми агрегацией. В данном конкретном примере вычисляется процент времени перелёта, который заняли задержки.

Предварительный просмотр показывает, что новый индекс будет содержать данные такого типа для каждой авиакомпании:

{
  "preview" : [
    {
      "carrier" : "ES-Air",
      "flights_count" : 2802.0,
      "flight_mins_total" : 1436927.5130677223,
      "delay_time_percentage" : 9.335543983955839,
      "delay_mins_total" : 134145.0
    },
    ...
  ]
}

Это преобразование упрощает ответы на такие вопросы, как:

  • Какая авиакомпания имеет наибольшее количество задержек в процентах от времени полета?

Эти данные являются вымышленными и не отражают фактические задержки или статистику полётов для каких-либо представленных аэропортов назначения или отправления.

Нахождение подозрительных IP-адресов клиентов

В этом примере используется набор образцовых данных веб-логов для выявления подозрительных IP-адресов клиентов. Он преобразует данные таким образом, что новый индекс содержит сумму байтов и количество уникальных URL-адресов, агентов, входящих запросов по местоположению и географическим пунктам назначения для каждого IP-адреса клиента. Он также использует агрегации фильтра для подсчета определённых типов HTTP-ответов, получаемых каждым IP-адресом клиента. В конечном итоге, пример ниже преобразует данные веб-логов в индекс, ориентированный на сущность, где сущностью является clientip.

resp = client.transform.put_transform(
    transform_id="suspicious_client_ips",
    source={
        "index": "kibana_sample_data_logs"
    },
    dest={
        "index": "sample_weblogs_by_clientip"
    },
    sync={
        "time": {
            "field": "timestamp",
            "delay": "60s"
        }
    },
    pivot={
        "group_by": {
            "clientip": {
                "terms": {
                    "field": "clientip"
                }
            }
        },
        "aggregations": {
            "url_dc": {
                "cardinality": {
                    "field": "url.keyword"
                }
            },
            "bytes_sum": {
                "sum": {
                    "field": "bytes"
                }
            },
            "geo.src_dc": {
                "cardinality": {
                    "field": "geo.src"
                }
            },
            "agent_dc": {
                "cardinality": {
                    "field": "agent.keyword"
                }
            },
            "geo.dest_dc": {
                "cardinality": {
                    "field": "geo.dest"
                }
            },
            "responses.total": {
                "value_count": {
                    "field": "timestamp"
                }
            },
            "success": {
                "filter": {
                    "term": {
                        "response": "200"
                    }
                }
            },
            "error404": {
                "filter": {
                    "term": {
                        "response": "404"
                    }
                }
            },
            "error5xx": {
                "filter": {
                    "range": {
                        "response": {
                            "gte": 500,
                            "lt": 600
                        }
                    }
                }
            },
            "timestamp.min": {
                "min": {
                    "field": "timestamp"
                }
            },
            "timestamp.max": {
                "max": {
                    "field": "timestamp"
                }
            },
            "timestamp.duration_ms": {
                "bucket_script": {
                    "buckets_path": {
                        "min_time": "timestamp.min.value",
                        "max_time": "timestamp.max.value"
                    },
                    "script": "(params.max_time - params.min_time)"
                }
            }
        }
    },
)
print(resp)
const response = await client.transform.putTransform({
  transform_id: "suspicious_client_ips",
  source: {
    index: "kibana_sample_data_logs",
  },
  dest: {
    index: "sample_weblogs_by_clientip",
  },
  sync: {
    time: {
      field: "timestamp",
      delay: "60s",
    },
  },
  pivot: {
    group_by: {
      clientip: {
        terms: {
          field: "clientip",
        },
      },
    },
    aggregations: {
      url_dc: {
        cardinality: {
          field: "url.keyword",
        },
      },
      bytes_sum: {
        sum: {
          field: "bytes",
        },
      },
      "geo.src_dc": {
        cardinality: {
          field: "geo.src",
        },
      },
      agent_dc: {
        cardinality: {
          field: "agent.keyword",
        },
      },
      "geo.dest_dc": {
        cardinality: {
          field: "geo.dest",
        },
      },
      "responses.total": {
        value_count: {
          field: "timestamp",
        },
      },
      success: {
        filter: {
          term: {
            response: "200",
          },
        },
      },
      error404: {
        filter: {
          term: {
            response: "404",
          },
        },
      },
      error5xx: {
        filter: {
          range: {
            response: {
              gte: 500,
              lt: 600,
            },
          },
        },
      },
      "timestamp.min": {
        min: {
          field: "timestamp",
        },
      },
      "timestamp.max": {
        max: {
          field: "timestamp",
        },
      },
      "timestamp.duration_ms": {
        bucket_script: {
          buckets_path: {
            min_time: "timestamp.min.value",
            max_time: "timestamp.max.value",
          },
          script: "(params.max_time - params.min_time)",
        },
      },
    },
  },
});
console.log(response);
PUT _transform/suspicious_client_ips
{
  "source": {
    "index": "kibana_sample_data_logs"
  },
  "dest" : { 
    "index" : "sample_weblogs_by_clientip"
  },
  "sync" : { 
    "time": {
      "field": "timestamp",
      "delay": "60s"
    }
  },
  "pivot": {
    "group_by": {  
      "clientip": { "terms": { "field": "clientip" } }
      },
    "aggregations": {
      "url_dc": { "cardinality": { "field": "url.keyword" }},
      "bytes_sum": { "sum": { "field": "bytes" }},
      "geo.src_dc": { "cardinality": { "field": "geo.src" }},
      "agent_dc": { "cardinality": { "field": "agent.keyword" }},
      "geo.dest_dc": { "cardinality": { "field": "geo.dest" }},
      "responses.total": { "value_count": { "field": "timestamp" }},
      "success" : { 
         "filter": {
            "term": { "response" : "200"}}
        },
      "error404" : {
         "filter": {
            "term": { "response" : "404"}}
        },
      "error5xx" : {
         "filter": {
            "range": { "response" : { "gte": 500, "lt": 600}}}
        },
      "timestamp.min": { "min": { "field": "timestamp" }},
      "timestamp.max": { "max": { "field": "timestamp" }},
      "timestamp.duration_ms": { 
        "bucket_script": {
          "buckets_path": {
            "min_time": "timestamp.min.value",
            "max_time": "timestamp.max.value"
          },
          "script": "(params.max_time - params.min_time)"
        }
      }
    }
  }
}

Целевой индекс для преобразования.

Настраивает преобразование для непрерывной работы. Он использует поле timestamp для синхронизации исходных и целевых индексов. Максимальная задержка обработки составляет 60 секунд.

Данные группируются по полю clientip.

Агрегация фильтра, подсчитывающая случаи успешных (200) ответов в поле response. Следующие две агрегации (error404 и error5xx) подсчитывают ответы с ошибками по кодам ошибок, совпадающим с точным значением или диапазоном кодов ответов.

Это bucket_script вычисляет продолжительность clientip доступа на основе результатов агрегации.

После создания преобразования необходимо запустить его:

resp = client.transform.start_transform(
    transform_id="suspicious_client_ips",
)
print(resp)
response = client.transform.start_transform(
  transform_id: 'suspicious_client_ips'
)
puts response
const response = await client.transform.startTransform({
  transform_id: "suspicious_client_ips",
});
console.log(response);
POST _transform/suspicious_client_ips/_start

Вскоре после этого первые результаты должны быть доступны в целевом индексе:

resp = client.search(
    index="sample_weblogs_by_clientip",
)
print(resp)
response = client.search(
  index: 'sample_weblogs_by_clientip'
)
puts response
const response = await client.search({
  index: "sample_weblogs_by_clientip",
});
console.log(response);
GET sample_weblogs_by_clientip/_search

Результат поиска показывает данные такого типа для каждого IP-адреса клиента:

    "hits" : [
      {
        "_index" : "sample_weblogs_by_clientip",
        "_id" : "MOeHH_cUL5urmartKj-b5UQAAAAAAAAA",
        "_score" : 1.0,
        "_source" : {
          "geo" : {
            "src_dc" : 2.0,
            "dest_dc" : 2.0
          },
          "success" : 2,
          "error404" : 0,
          "error503" : 0,
          "clientip" : "0.72.176.46",
          "agent_dc" : 2.0,
          "bytes_sum" : 4422.0,
          "responses" : {
            "total" : 2.0
          },
          "url_dc" : 2.0,
          "timestamp" : {
            "duration_ms" : 5.2191698E8,
            "min" : "2020-03-16T07:51:57.333Z",
            "max" : "2020-03-22T08:50:34.313Z"
          }
        }
      }
    ]

Как и другие наборы образцовых данных Kibana, набор образцовых данных веб-логов содержит временные метки, относительные к моменту установки, включая временные метки в будущем. Непрерывное преобразование будет собирать данные, когда они окажутся в прошлом. Если вы установили набор образцовых данных веб-логов некоторое время назад, вы можете его удалить и переустановить, и временные метки изменятся.

Это преобразование упрощает ответы на такие вопросы, как:

  • Какие IP-адреса клиентов передают наибольшее количество данных?
  • Какие IP-адреса клиентов взаимодействуют с большим количеством различных URL-адресов?
  • Какие IP-адреса клиентов имеют высокий показатель ошибок?
  • Какие IP-адреса клиентов взаимодействуют с большим количеством стран назначения?

Поиск последнего события регистрации для каждого IP-адреса

В этом примере используется набор данных примеров веб-логов для поиска последнего лога с определенного IP-адреса. Воспользуемся трансформером типа latest в непрерывном режиме. Он копирует самый последний документ для каждого уникального ключа из исходного индекса в целевой индекс и обновляет целевой индекс по мере поступления новых данных в исходный индекс.

Выберите поле clientip в качестве уникального ключа; данные группируются по этому полю. Выберите поле timestamp в качестве поля даты, которое сортирует данные в хронологическом порядке. Для непрерывного режима укажите поле даты, используемое для идентификации новых документов, и интервал между проверками изменений в исходном индексе.

Finding the last log event for each IP address with transforms in Kibana

Предположим, нас интересует сохранение документов только для IP-адресов, которые недавно появились в логах. Вы можете определить политику удержания и указать поле даты, которое используется для расчета возраста документа. В этом примере используется то же поле даты, которое используется для сортировки данных. Затем установите максимальный возраст документа; документы, старше указанного значения, будут удалены из целевого индекса.

Defining retention policy for transforms in Kibana

Этот трансформер создаёт целевой индекс, содержащий последнюю дату входа для каждого клиентского IP-адреса. Поскольку трансформер работает в непрерывном режиме, целевой индекс будет обновляться по мере поступления новых данных в исходный индекс. Наконец, каждый документ, старше 30 дней, будет удалён из целевого индекса из-за применённой политики удержания.

Пример API
resp = client.transform.put_transform(
    transform_id="last-log-from-clientip",
    source={
        "index": [
            "kibana_sample_data_logs"
        ]
    },
    latest={
        "unique_key": [
            "clientip"
        ],
        "sort": "timestamp"
    },
    frequency="1m",
    dest={
        "index": "last-log-from-clientip"
    },
    sync={
        "time": {
            "field": "timestamp",
            "delay": "60s"
        }
    },
    retention_policy={
        "time": {
            "field": "timestamp",
            "max_age": "30d"
        }
    },
    settings={
        "max_page_search_size": 500
    },
)
print(resp)
const response = await client.transform.putTransform({
  transform_id: "last-log-from-clientip",
  source: {
    index: ["kibana_sample_data_logs"],
  },
  latest: {
    unique_key: ["clientip"],
    sort: "timestamp",
  },
  frequency: "1m",
  dest: {
    index: "last-log-from-clientip",
  },
  sync: {
    time: {
      field: "timestamp",
      delay: "60s",
    },
  },
  retention_policy: {
    time: {
      field: "timestamp",
      max_age: "30d",
    },
  },
  settings: {
    max_page_search_size: 500,
  },
});
console.log(response);
PUT _transform/last-log-from-clientip
{
  "source": {
    "index": [
      "kibana_sample_data_logs"
    ]
  },
  "latest": {
    "unique_key": [ 
      "clientip"
    ],
    "sort": "timestamp" 
  },
  "frequency": "1m", 
  "dest": {
    "index": "last-log-from-clientip"
  },
  "sync": { 
    "time": {
      "field": "timestamp",
      "delay": "60s"
    }
  },
  "retention_policy": { 
    "time": {
      "field": "timestamp",
      "max_age": "30d"
    }
  },
  "settings": {
    "max_page_search_size": 500
  }
}

Указывает поле для группировки данных.

Указывает поле даты, используемое для сортировки данных.

Устанавливает интервал для трансформера, чтобы проверять изменения в исходном индексе.

Содержит поле времени и настройки задержки, используемые для синхронизации исходного и целевого индексов.

Указывает политику удержания для трансформера. Документы, старше заданного значения, будут удалены из целевого индекса.

После создания трансформера запустите его:

resp = client.transform.start_transform(
    transform_id="last-log-from-clientip",
)
print(resp)
response = client.transform.start_transform(
  transform_id: 'last-log-from-clientip'
)
puts response
const response = await client.transform.startTransform({
  transform_id: "last-log-from-clientip",
});
console.log(response);
POST _transform/last-log-from-clientip/_start

После обработки данных трансформером выполните поиск в целевом индексе:

resp = client.search(
    index="last-log-from-clientip",
)
print(resp)
response = client.search(
  index: 'last-log-from-clientip'
)
puts response
const response = await client.search({
  index: "last-log-from-clientip",
});
console.log(response);
GET last-log-from-clientip/_search

Результат поиска покажет данные такого типа для каждого клиентского IP:

{
  "_index" : "last-log-from-clientip",
  "_id" : "MOeHH_cUL5urmartKj-b5UQAAAAAAAAA",
  "_score" : 1.0,
  "_source" : {
    "referer" : "http://twitter.com/error/don-lind",
    "request" : "/elasticsearch",
    "agent" : "Mozilla/4.0 (compatible; MSIE 6.0; Windows NT 5.1; SV1; .NET CLR 1.1.4322)",
    "extension" : "",
    "memory" : null,
    "ip" : "0.72.176.46",
    "index" : "kibana_sample_data_logs",
    "message" : "0.72.176.46 - - [2018-09-18T06:31:00.572Z] \"GET /elasticsearch HTTP/1.1\" 200 7065 \"-\" \"Mozilla/4.0 (compatible; MSIE 6.0; Windows NT 5.1; SV1; .NET CLR 1.1.4322)\"",
    "url" : "https://www.elastic.co/downloads/elasticsearch",
    "tags" : [
      "success",
      "info"
    ],
    "geo" : {
      "srcdest" : "IN:PH",
      "src" : "IN",
      "coordinates" : {
        "lon" : -124.1127917,
        "lat" : 40.80338889
      },
      "dest" : "PH"
    },
    "utc_time" : "2021-05-04T06:31:00.572Z",
    "bytes" : 7065,
    "machine" : {
      "os" : "ios",
      "ram" : 12884901888
    },
    "response" : 200,
    "clientip" : "0.72.176.46",
    "host" : "www.elastic.co",
    "event" : {
      "dataset" : "sample_web_logs"
    },
    "phpmemory" : null,
    "timestamp" : "2021-05-04T06:31:00.572Z"
  }
}

Этот трансформер упрощает ответы на вопросы, такие как:

  • Какое было последнее событие регистрации, связанное с конкретным IP-адресом?

Поиск IP-адресов клиентов, которые отправили больше всего байтов на сервер

В этом примере используется набор данных примеров веб-логов для поиска IP-адреса клиента, который отправил больше всего байтов на сервер за каждый час. В примере используется трансформер pivot с агрегацией top_metrics.

Сгруппируйте данные по гистограмме дат по полю времени с интервалом в один час. Используйте максимальную агрегацию по полю bytes, чтобы получить максимальное количество данных, отправленных на сервер. Без агрегации max API-вызов всё равно возвращает IP-адрес клиента, который отправил больше всего байтов, однако количество отправленных байтов не возвращается. В свойстве top_metrics укажите clientip и geo.src, затем отсортируйте их по полю bytes в порядке убывания. Трансформер возвращает IP-адрес клиента, который отправил наибольшее количество данных, и двухбуквенный ISO-код соответствующего местоположения.

resp = client.transform.preview_transform(
    source={
        "index": "kibana_sample_data_logs"
    },
    pivot={
        "group_by": {
            "timestamp": {
                "date_histogram": {
                    "field": "timestamp",
                    "fixed_interval": "1h"
                }
            }
        },
        "aggregations": {
            "bytes.max": {
                "max": {
                    "field": "bytes"
                }
            },
            "top": {
                "top_metrics": {
                    "metrics": [
                        {
                            "field": "clientip"
                        },
                        {
                            "field": "geo.src"
                        }
                    ],
                    "sort": {
                        "bytes": "desc"
                    }
                }
            }
        }
    },
)
print(resp)
const response = await client.transform.previewTransform({
  source: {
    index: "kibana_sample_data_logs",
  },
  pivot: {
    group_by: {
      timestamp: {
        date_histogram: {
          field: "timestamp",
          fixed_interval: "1h",
        },
      },
    },
    aggregations: {
      "bytes.max": {
        max: {
          field: "bytes",
        },
      },
      top: {
        top_metrics: {
          metrics: [
            {
              field: "clientip",
            },
            {
              field: "geo.src",
            },
          ],
          sort: {
            bytes: "desc",
          },
        },
      },
    },
  },
});
console.log(response);
POST _transform/_preview
{
  "source": {
    "index": "kibana_sample_data_logs"
  },
  "pivot": {
    "group_by": { 
      "timestamp": {
        "date_histogram": {
          "field": "timestamp",
          "fixed_interval": "1h"
        }
      }
    },
    "aggregations": {
      "bytes.max": { 
        "max": {
          "field": "bytes"
        }
      },
      "top": {
        "top_metrics": { 
          "metrics": [
            {
              "field": "clientip"
            },
            {
              "field": "geo.src"
            }
          ],
          "sort": {
            "bytes": "desc"
          }
        }
      }
    }
  }
}

Данные группируются по гистограмме дат по полю времени с интервалом в один час.

Вычисляется максимальное значение поля bytes.

Указываются поля (clientip и geo.src) верхнего документа для возврата и способ сортировки (документ с наибольшим значением поля bytes).

Вышеприведённый API-вызов возвращает ответ, похожий на этот:

{
  "preview" : [
    {
      "top" : {
        "clientip" : "223.87.60.27",
        "geo.src" : "IN"
      },
      "bytes" : {
        "max" : 6219
      },
      "timestamp" : "2021-04-25T00:00:00.000Z"
    },
    {
      "top" : {
        "clientip" : "99.74.118.237",
        "geo.src" : "LK"
      },
      "bytes" : {
        "max" : 14113
      },
      "timestamp" : "2021-04-25T03:00:00.000Z"
    },
    {
      "top" : {
        "clientip" : "218.148.135.12",
        "geo.src" : "BR"
      },
      "bytes" : {
        "max" : 4531
      },
      "timestamp" : "2021-04-25T04:00:00.000Z"
    },
    ...
  ]
}

Получение имени и адреса электронной почты клиента по идентификатору клиента

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

Сгруппируйте данные по customer_id, затем добавьте агрегацию top_metrics, где поля metrics являются полями email, customer_first_name.keyword и customer_last_name.keyword. Отсортируйте top_metrics по полю order_date в порядке убывания. API-вызов выглядит следующим образом:

resp = client.transform.preview_transform(
    source={
        "index": "kibana_sample_data_ecommerce"
    },
    pivot={
        "group_by": {
            "customer_id": {
                "terms": {
                    "field": "customer_id"
                }
            }
        },
        "aggregations": {
            "last": {
                "top_metrics": {
                    "metrics": [
                        {
                            "field": "email"
                        },
                        {
                            "field": "customer_first_name.keyword"
                        },
                        {
                            "field": "customer_last_name.keyword"
                        }
                    ],
                    "sort": {
                        "order_date": "desc"
                    }
                }
            }
        }
    },
)
print(resp)
const response = await client.transform.previewTransform({
  source: {
    index: "kibana_sample_data_ecommerce",
  },
  pivot: {
    group_by: {
      customer_id: {
        terms: {
          field: "customer_id",
        },
      },
    },
    aggregations: {
      last: {
        top_metrics: {
          metrics: [
            {
              field: "email",
            },
            {
              field: "customer_first_name.keyword",
            },
            {
              field: "customer_last_name.keyword",
            },
          ],
          sort: {
            order_date: "desc",
          },
        },
      },
    },
  },
});
console.log(response);
POST _transform/_preview
{
  "source": {
    "index": "kibana_sample_data_ecommerce"
  },
  "pivot": {
    "group_by": { 
      "customer_id": {
        "terms": {
          "field": "customer_id"
        }
      }
    },
    "aggregations": {
      "last": {
        "top_metrics": { 
          "metrics": [
            {
              "field": "email"
            },
            {
              "field": "customer_first_name.keyword"
            },
            {
              "field": "customer_last_name.keyword"
            }
          ],
          "sort": {
            "order_date": "desc"
          }
        }
      }
    }
  }
}

Данные группируются с помощью агрегации terms по полю customer_id.

Указываются поля для возврата (поля email и name) в порядке убывания по полю order date.

API возвращает ответ, похожий на этот:

 {
  "preview" : [
    {
      "last" : {
        "customer_last_name.keyword" : "Long",
        "customer_first_name.keyword" : "Recip",
        "email" : "recip@long-family.zzz"
      },
      "customer_id" : "10"
    },
    {
      "last" : {
        "customer_last_name.keyword" : "Jackson",
        "customer_first_name.keyword" : "Fitzgerald",
        "email" : "fitzgerald@jackson-family.zzz"
      },
      "customer_id" : "11"
    },
    {
      "last" : {
        "customer_last_name.keyword" : "Cross",
        "customer_first_name.keyword" : "Brigitte",
        "email" : "brigitte@cross-family.zzz"
      },
      "customer_id" : "12"
    },
    ...
  ]
}

© 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/transform-examples.html

Spec-Zone.ru

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