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

Примеры использования Painless для преобразований

Примеры, использующие агрегацию scripted_metric, не поддерживаются в Elasticsearch Serverless.

Эти примеры демонстрируют использование Painless в преобразованиях. Дополнительную информацию о языке сценариев Painless можно найти в руководстве по Painless.

  • Получение лучших совпадений с помощью агрегации списанных метрик
  • Получение временных характеристик с помощью агрегаций
  • Получение продолжительности с помощью скрипта для ведер
  • Подсчет HTTP-ответов с помощью агрегации списанных метрик
  • Сравнение индексов с помощью агрегаций списанных метрик
  • Получение подробностей о веб-сессиях с помощью агрегации списанных метрик
  • Хотя контекст следующих примеров — преобразование, скрипты Painless в приведенных ниже фрагментах кода также можно использовать в других агрегациях поиска Elasticsearch.
  • Все следующие примеры используют скрипты; преобразования не могут вывести сопоставления выходных полей, когда поля создаются скриптом. Преобразования не создают никаких сопоставлений в целевом индексе для этих полей, что означает, что они сопоставляются динамически. Создайте целевой индекс перед запуском преобразования, если вы хотите явные сопоставления.

Получение лучших совпадений с помощью агрегации списанных метрик

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

В этом примере используется агрегация scripted_metric, которая не поддерживается в Elasticsearch Serverless.

"aggregations": {
  "latest_doc": {
    "scripted_metric": {
      "init_script": "state.timestamp_latest = 0L; state.last_doc = ''", 
      "map_script": """ 
        def current_date = doc['@timestamp'].getValue().toInstant().toEpochMilli();
        if (current_date > state.timestamp_latest)
        {state.timestamp_latest = current_date;
        state.last_doc = new HashMap(params['_source']);}
      """,
      "combine_script": "return state", 
      "reduce_script": """ 
        def last_doc = '';
        def timestamp_latest = 0L;
        for (s in states) {if (s.timestamp_latest > (timestamp_latest))
        {timestamp_latest = s.timestamp_latest; last_doc = s.last_doc;}}
        return last_doc
      """
    }
  }
}

init_script создает поле типа long timestamp_latest и строкового типа last_doc в объекте state.

map_script определяет current_date на основе метки времени документа, затем сравнивает current_date с state.timestamp_latest и, наконец, возвращает state.last_doc из фрагмента. Используя new HashMap(...), вы копируете исходный документ, это важно, когда вам нужно передать весь исходный объект с одной фазы на другую.

combine_script возвращает state из каждого фрагмента.

reduce_script итерируется по значению s.timestamp_latest, возвращаемому каждым фрагментом, и возвращает документ с самой поздней меткой времени (last_doc). В ответ лучшее совпадение (другими словами, latest_doc) вложено в поле latest_doc.

Подробное описание области действия скриптов см. в разделе область действия скриптов.

Вы можете получить последнее значение аналогичным способом:

"aggregations": {
  "latest_value": {
    "scripted_metric": {
      "init_script": "state.timestamp_latest = 0L; state.last_value = ''",
      "map_script": """
        def current_date = doc['@timestamp'].getValue().toInstant().toEpochMilli();
        if (current_date > state.timestamp_latest)
        {state.timestamp_latest = current_date;
        state.last_value = params['_source']['value'];}
      """,
      "combine_script": "return state",
      "reduce_script": """
        def last_value = '';
        def timestamp_latest = 0L;
        for (s in states) {if (s.timestamp_latest > (timestamp_latest))
        {timestamp_latest = s.timestamp_latest; last_value = s.last_value;}}
        return last_value
      """
    }
  }
}
Получение лучших совпадений с помощью хранимых скриптов

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

  1. Создайте хранимые скрипты:

    POST _scripts/last-value-map-init
    {
      "script": {
        "lang": "painless",
        "source": """
            state.timestamp_latest = 0L; state.last_value = ''
        """
      }
    }
    
    POST _scripts/last-value-map
    {
      "script": {
        "lang": "painless",
        "source": """
          def current_date = doc['@timestamp'].getValue().toInstant().toEpochMilli();
            if (current_date > state.timestamp_latest)
            {state.timestamp_latest = current_date;
            state.last_value = doc[params['key']].value;}
        """
      }
    }
    
    POST _scripts/last-value-combine
    {
      "script": {
        "lang": "painless",
        "source": """
            return state
        """
      }
    }
    
    POST _scripts/last-value-reduce
    {
      "script": {
        "lang": "painless",
        "source": """
            def last_value = '';
            def timestamp_latest = 0L;
            for (s in states) {if (s.timestamp_latest > (timestamp_latest))
            {timestamp_latest = s.timestamp_latest; last_value = s.last_value;}}
            return last_value
        """
      }
    }
  2. Используйте хранимые скрипты в агрегации списанных метрик.

    "aggregations":{
       "latest_value":{
          "scripted_metric":{
             "init_script":{
                "id":"last-value-map-init"
             },
             "map_script":{
                "id":"last-value-map",
                "params":{
                   "key":"field_with_last_value" 
                }
             },
             "combine_script":{
                "id":"last-value-combine"
             },
             "reduce_script":{
                "id":"last-value-reduce"
             }

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

Получение временных характеристик с помощью агрегаций

Этот фрагмент демонстрирует, как извлекать временные характеристики, используя Painless в преобразовании. Фрагмент использует индекс, в котором @timestamp определен как поле типа date.

"aggregations": {
  "avg_hour_of_day": { 
    "avg":{
      "script": { 
        "source": """
          ZonedDateTime date =  doc['@timestamp'].value; 
          return date.getHour(); 
        """
      }
    }
  },
  "avg_month_of_year": { 
    "avg":{
      "script": { 
        "source": """
          ZonedDateTime date =  doc['@timestamp'].value; 
          return date.getMonthValue(); 
        """
      }
    }
  },
 ...
}

Имя агрегации.

Содержит скрипт Painless, возвращающий час дня.

Устанавливает date на основе метки времени документа.

Возвращает значение часа из date.

Имя агрегации.

Содержит скрипт Painless, возвращающий месяц года.

Устанавливает date на основе метки времени документа.

Возвращает значение месяца из date.

Получение продолжительности с помощью скрипта для ведер

В этом примере показано, как получить продолжительность сессии по IP-адресу клиента из журнала данных, используя bucket script. В примере используется набор данных веб-журналов Kibana.

resp = client.transform.put_transform(
    transform_id="data_log",
    source={
        "index": "kibana_sample_data_logs"
    },
    dest={
        "index": "data-logs-by-client"
    },
    pivot={
        "group_by": {
            "machine.os": {
                "terms": {
                    "field": "machine.os.keyword"
                }
            },
            "machine.ip": {
                "terms": {
                    "field": "clientip"
                }
            }
        },
        "aggregations": {
            "time_frame.lte": {
                "max": {
                    "field": "timestamp"
                }
            },
            "time_frame.gte": {
                "min": {
                    "field": "timestamp"
                }
            },
            "time_length": {
                "bucket_script": {
                    "buckets_path": {
                        "min": "time_frame.gte.value",
                        "max": "time_frame.lte.value"
                    },
                    "script": "params.max - params.min"
                }
            }
        }
    },
)
print(resp)
const response = await client.transform.putTransform({
  transform_id: "data_log",
  source: {
    index: "kibana_sample_data_logs",
  },
  dest: {
    index: "data-logs-by-client",
  },
  pivot: {
    group_by: {
      "machine.os": {
        terms: {
          field: "machine.os.keyword",
        },
      },
      "machine.ip": {
        terms: {
          field: "clientip",
        },
      },
    },
    aggregations: {
      "time_frame.lte": {
        max: {
          field: "timestamp",
        },
      },
      "time_frame.gte": {
        min: {
          field: "timestamp",
        },
      },
      time_length: {
        bucket_script: {
          buckets_path: {
            min: "time_frame.gte.value",
            max: "time_frame.lte.value",
          },
          script: "params.max - params.min",
        },
      },
    },
  },
});
console.log(response);
PUT _transform/data_log
{
  "source": {
    "index": "kibana_sample_data_logs"
  },
  "dest": {
    "index": "data-logs-by-client"
  },
  "pivot": {
    "group_by": {
      "machine.os": {"terms": {"field": "machine.os.keyword"}},
      "machine.ip": {"terms": {"field": "clientip"}}
    },
    "aggregations": {
      "time_frame.lte": {
        "max": {
          "field": "timestamp"
        }
      },
      "time_frame.gte": {
        "min": {
          "field": "timestamp"
        }
      },
      "time_length": { 
        "bucket_script": {
          "buckets_path": { 
            "min": "time_frame.gte.value",
            "max": "time_frame.lte.value"
          },
          "script": "params.max - params.min" 
        }
      }
    }
  }
}

Для определения длительности сессий мы используем скрипт для ведер.

Путь ведра — это карта переменных скрипта и соответствующего пути к ведрам, которые вы хотите использовать для переменной. В данном случае, min и max сопоставлены с time_frame.gte.value и time_frame.lte.value.

Наконец, скрипт вычитает начальную дату сессии из конечной даты, что приводит к продолжительности сессии.

Подсчет HTTP-ответов с помощью скриптовой агрегации метрик

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

В примере ниже предполагается, что коды HTTP-ответов хранятся в качестве ключевых слов в поле response документов.

В этом примере используется агрегация scripted_metric, которая не поддерживается в Elasticsearch Serverless.

"aggregations": { 
  "responses.counts": { 
    "scripted_metric": { 
      "init_script": "state.responses = ['error':0L,'success':0L,'other':0L]", 
      "map_script": """ 
        def code = doc['response.keyword'].value;
        if (code.startsWith('5') || code.startsWith('4')) {
          state.responses.error += 1 ;
        } else if(code.startsWith('2')) {
          state.responses.success += 1;
        } else {
          state.responses.other += 1;
        }
        """,
      "combine_script": "state.responses", 
      "reduce_script": """ 
        def counts = ['error': 0L, 'success': 0L, 'other': 0L];
        for (responses in states) {
          counts.error += responses['error'];
          counts.success += responses['success'];
          counts.other += responses['other'];
        }
        return counts;
        """
      }
    },
  ...
}

Объект aggregations трансформации, содержащий все агрегации.

Объект агрегации scripted_metric.

Данная scripted_metric выполняет распределенную операцию над данными веб-логов для подсчета определенных типов HTTP-ответов (ошибки, успехи и другие).

init_script создает массив responses в объекте state с тремя свойствами (error, success, other) со значением типа long.

map_script определяет code на основе значения response.keyword документа, затем подсчитывает ошибки, успехи и другие ответы на основе первой цифры кода ответа.

combine_script возвращает state.responses с каждого фрагмента.

reduce_script создаёт массив counts со свойствами error, success и other, затем итерируется по значению responses, возвращаемому каждым фрагментом, и присваивает различные типы ответов соответствующим свойствам объекта counts; ответы об ошибках — счётчикам ошибок, успешные ответы — счётчикам успехов, а другие ответы — счётчикам других ответов. Наконец, возвращает массив counts со счётчиками ответов.

Сравнение индексов с помощью скриптовой агрегации метрик

Этот пример демонстрирует, как сравнить содержимое двух индексов с помощью трансформации, использующей скриптовую агрегацию метрик.

В этом примере используется агрегация scripted_metric, которая не поддерживается в Elasticsearch Serverless.

resp = client.transform.preview_transform(
    id="index_compare",
    source={
        "index": [
            "index1",
            "index2"
        ],
        "query": {
            "match_all": {}
        }
    },
    dest={
        "index": "compare"
    },
    pivot={
        "group_by": {
            "unique-id": {
                "terms": {
                    "field": "<unique-id-field>"
                }
            }
        },
        "aggregations": {
            "compare": {
                "scripted_metric": {
                    "map_script": "state.doc = new HashMap(params['_source'])",
                    "combine_script": "return state",
                    "reduce_script": " \n            if (states.size() != 2) {\n              return \"count_mismatch\"\n            }\n            if (states.get(0).equals(states.get(1))) {\n              return \"match\"\n            } else {\n              return \"mismatch\"\n            }\n            "
                }
            }
        }
    },
)
print(resp)
const response = await client.transform.previewTransform({
  id: "index_compare",
  source: {
    index: ["index1", "index2"],
    query: {
      match_all: {},
    },
  },
  dest: {
    index: "compare",
  },
  pivot: {
    group_by: {
      "unique-id": {
        terms: {
          field: "<unique-id-field>",
        },
      },
    },
    aggregations: {
      compare: {
        scripted_metric: {
          map_script: "state.doc = new HashMap(params['_source'])",
          combine_script: "return state",
          reduce_script:
            ' \n            if (states.size() != 2) {\n              return "count_mismatch"\n            }\n            if (states.get(0).equals(states.get(1))) {\n              return "match"\n            } else {\n              return "mismatch"\n            }\n            ',
        },
      },
    },
  },
});
console.log(response);
POST _transform/_preview
{
  "id" : "index_compare",
  "source" : { 
    "index" : [
      "index1",
      "index2"
    ],
    "query" : {
      "match_all" : { }
    }
  },
  "dest" : { 
    "index" : "compare"
  },
  "pivot" : {
    "group_by" : {
      "unique-id" : {
        "terms" : {
          "field" : "<unique-id-field>" 
        }
      }
    },
    "aggregations" : {
      "compare" : { 
        "scripted_metric" : {
          "map_script" : "state.doc = new HashMap(params['_source'])", 
          "combine_script" : "return state", 
          "reduce_script" : """ 
            if (states.size() != 2) {
              return "count_mismatch"
            }
            if (states.get(0).equals(states.get(1))) {
              return "match"
            } else {
              return "mismatch"
            }
            """
        }
      }
    }
  }
}

Индексы, указанные в объекте source, сравниваются друг с другом.

Индекс dest содержит результаты сравнения.

Поле group_by должно быть уникальным идентификатором каждого документа.

Объект агрегации scripted_metric.

map_script определяет doc в объекте состояния. Использование new HashMap(...) позволяет скопировать исходный документ, что важно, когда необходимо передать весь исходный объект из одной фазы в следующую.

combine_script возвращает state с каждого фрагмента.

reduce_script проверяет, равны ли размеры индексов. Если они не равны, то возвращает сообщение count_mismatch. Затем он итерируется по всем значениям двух индексов и сравнивает их. Если значения равны, возвращается match, в противном случае — mismatch.

Получение деталей сеанса веб-пользователя с помощью скриптовой агрегации метрик

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

Исходный документ
{
  "_index":"apache-sessions",
  "_type":"_doc",
  "_id":"KvzSeGoB4bgw0KGbE3wP",
  "_score":1.0,
  "_source":{
    "@timestamp":1484053499256,
    "apache":{
      "access":{
        "sessionid":"571604f2b2b0c7b346dc685eeb0e2306774a63c2",
        "url":"http://www.leroymerlin.fr/v3/search/search.do?keyword=Carrelage%20salle%20de%20bain",
        "path":"/v3/search/search.do",
        "query":"keyword=Carrelage%20salle%20de%20bain",
        "referrer":"http://www.leroymerlin.fr/v3/p/produits/carrelage-parquet-sol-souple/carrelage-sol-et-mur/decor-listel-et-accessoires-carrelage-mural-l1308217717?resultOffset=0&resultLimit=51&resultListShape=MOSAIC&priceStyle=SALEUNIT_PRICE",
        "user_agent":{
          "original":"Mobile Safari 10.0 Mac OS X (iPad) Apple Inc.",
          "os_name":"Mac OS X (iPad)"
        },
        "remote_ip":"0337b1fa-5ed4-af81-9ef4-0ec53be0f45d",
        "geoip":{
          "country_iso_code":"FR",
          "location":{
            "lat":48.86,
            "lon":2.35
          }
        },
        "response_code":200,
        "method":"GET"
      }
    }
  }
}
...

Использование sessionid в качестве поля группировки позволяет перечислить события в сеансе и получить более подробную информацию о нём с помощью скриптовой агрегации метрик.

В этом примере используется агрегация scripted_metric, которая не поддерживается в Elasticsearch Serverless.

POST _transform/_preview
{
  "source": {
    "index": "apache-sessions"
  },
  "pivot": {
    "group_by": {
      "sessionid": { 
        "terms": {
          "field": "apache.access.sessionid"
        }
      }
    },
    "aggregations": { 
      "distinct_paths": {
        "cardinality": {
          "field": "apache.access.path"
        }
      },
      "num_pages_viewed": {
        "value_count": {
          "field": "apache.access.url"
        }
      },
      "session_details": {
        "scripted_metric": {
          "init_script": "state.docs = []", 
          "map_script": """ 
            Map span = [
              '@timestamp':doc['@timestamp'].value,
              'url':doc['apache.access.url'].value,
              'referrer':doc['apache.access.referrer'].value
            ];
            state.docs.add(span)
          """,
          "combine_script": "return state.docs;", 
          "reduce_script": """ 
            def all_docs = [];
            for (s in states) {
              for (span in s) {
                all_docs.add(span);
              }
            }
            all_docs.sort((HashMap o1, HashMap o2)->o1['@timestamp'].toEpochMilli().compareTo(o2['@timestamp'].toEpochMilli()));
            def size = all_docs.size();
            def min_time = all_docs[0]['@timestamp'];
            def max_time = all_docs[size-1]['@timestamp'];
            def duration = max_time.toEpochMilli() - min_time.toEpochMilli();
            def entry_page = all_docs[0]['url'];
            def exit_path = all_docs[size-1]['url'];
            def first_referrer = all_docs[0]['referrer'];
            def ret = new HashMap();
            ret['first_time'] = min_time;
            ret['last_time'] = max_time;
            ret['duration'] = duration;
            ret['entry_page'] = entry_page;
            ret['exit_path'] = exit_path;
            ret['first_referrer'] = first_referrer;
            return ret;
          """
        }
      }
    }
  }
}

Данные сгруппированы по полю sessionid.

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

init_script создаёт массив типа doc в объекте state.

map_script определяет массив span со временем, URL и значением referrer, которые основаны на соответствующих значениях документа, а затем добавляет значение массива span в объект doc.

combine_script возвращает state.docs с каждого фрагмента.

reduce_script определяет различные объекты, такие как min_time, max_time и duration, на основе полей документа, затем объявляет объект ret и копирует исходный документ с помощью new HashMap (). Затем скрипт определяет поля first_time, last_time, duration и другие в объекте ret на основе ранее определённых объектов, и наконец, возвращает ret.

Результат вызова API похож на этот:

{
  "num_pages_viewed" : 2.0,
  "session_details" : {
    "duration" : 100300001,
    "first_referrer" : "https://www.bing.com/",
    "entry_page" : "http://www.leroymerlin.fr/v3/p/produits/materiaux-menuiserie/porte-coulissante-porte-interieure-escalier-et-rambarde/barriere-de-securite-l1308218463",
    "first_time" : "2017-01-10T21:22:52.982Z",
    "last_time" : "2017-01-10T21:25:04.356Z",
    "exit_path" : "http://www.leroymerlin.fr/v3/p/produits/materiaux-menuiserie/porte-coulissante-porte-interieure-escalier-et-rambarde/barriere-de-securite-l1308218463?__result-wrapper?pageTemplate=Famille%2FMat%C3%A9riaux+et+menuiserie&resultOffset=0&resultLimit=50&resultListShape=PLAIN&nomenclatureId=17942&priceStyle=SALEUNIT_PRICE&fcr=1&*4294718806=4294718806&*14072=14072&*4294718593=4294718593&*17942=17942"
  },
  "distinct_paths" : 1.0,
  "sessionid" : "000046f8154a80fd89849369c984b8cc9d795814"
},
{
  "num_pages_viewed" : 10.0,
  "session_details" : {
    "duration" : 343100405,
    "first_referrer" : "https://www.google.fr/",
    "entry_page" : "http://www.leroymerlin.fr/",
    "first_time" : "2017-01-10T16:57:39.937Z",
    "last_time" : "2017-01-10T17:03:23.049Z",
    "exit_path" : "http://www.leroymerlin.fr/v3/p/produits/porte-de-douche-coulissante-adena-e168578"
  },
  "distinct_paths" : 8.0,
  "sessionid" : "000087e825da1d87a332b8f15fa76116c7467da6"
}
...

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

Spec-Zone.ru

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