Spec-Zone.ru › Elasticsearch 7
›Руководство по Elasticsearch [7.17]

Конвейеры Ingest

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

Конвейер состоит из ряда настраиваемых задач, называемых процессорами. Каждый процессор выполняется последовательно, внося определённые изменения в поступающие документы. После выполнения процессоров Elasticsearch добавляет преобразованные документы в ваш поток данных или индекс.

Ingest pipeline diagram

Вы можете создавать и управлять конвейерами загрузки, используя функцию Конвейеры загрузки Kibana или API загрузки. Elasticsearch хранит конвейеры в состоянии кластера.

Предварительные требования

  • Узлы с ролью узла ingest обрабатывают обработку конвейера. Для использования конвейеров загрузки в вашем кластере должен быть как минимум один узел с ролью ingest. Для больших нагрузок загрузки рекомендуется создавать посвящённые узлы загрузки.
  • Если включены функции безопасности Elasticsearch, вам необходимо иметь право manage_pipeline доступа к кластеру для управления конвейерами загрузки. Для использования функции Конвейеры загрузки Kibana вам также потребуются права cluster:monitor/nodes/info доступа к кластеру.
  • Конвейеры, включающие процессор enrich, требуют дополнительной настройки. См. Обогащение ваших данных.

Создание и управление конвейерами

В Kibana откройте главное меню и нажмите Управление стеком > Конвейеры загрузки. Из представления списка вы можете:

  • Просмотреть список ваших конвейеров и углубиться в детали
  • Изменять или клонировать существующие конвейеры
  • Удалять конвейеры

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

Kibana’s Ingest Pipelines list view

Вы также можете использовать API загрузки для создания и управления конвейерами. Следующий запрос API создания конвейера создает конвейер, содержащий два set процессора, за которыми следует lowercase процессор. Процессоры выполняются последовательно в указанном порядке.

PUT _ingest/pipeline/my-pipeline
{
  "description": "My optional pipeline description",
  "processors": [
    {
      "set": {
        "description": "My optional processor description",
        "field": "my-long-field",
        "value": 10
      }
    },
    {
      "set": {
        "description": "Set 'my-boolean-field' to true",
        "field": "my-boolean-field",
        "value": true
      }
    },
    {
      "lowercase": {
        "field": "my-keyword-field"
      }
    }
  ]
}

Управление версиями конвейера

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

PUT _ingest/pipeline/my-pipeline-id
{
  "version": 1,
  "processors": [ ... ]
}

Чтобы сбросить version число с помощью API, замените или обновите конвейер, не указывая параметр version.

Тестирование конвейера

Перед использованием конвейера в рабочей среде рекомендуется протестировать его с помощью образцов документов. При создании или редактировании конвейера в Kibana нажмите Добавить документы. В вкладке Документы укажите образцы документов и нажмите Запустить конвейер.

Test a pipeline in Kibana

Вы также можете протестировать конвейеры с помощью API моделирования конвейера. Вы можете указать настроенный конвейер в пути запроса. Например, следующий запрос тестирует my-pipeline.

POST _ingest/pipeline/my-pipeline/_simulate
{
  "docs": [
    {
      "_source": {
        "my-keyword-field": "FOO"
      }
    },
    {
      "_source": {
        "my-keyword-field": "BAR"
      }
    }
  ]
}

В качестве альтернативы вы можете указать конвейер и его процессоры в теле запроса.

POST _ingest/pipeline/_simulate
{
  "pipeline": {
    "processors": [
      {
        "lowercase": {
          "field": "my-keyword-field"
        }
      }
    ]
  },
  "docs": [
    {
      "_source": {
        "my-keyword-field": "FOO"
      }
    },
    {
      "_source": {
        "my-keyword-field": "BAR"
      }
    }
  ]
}

API возвращает преобразованные документы:

{
  "docs": [
    {
      "doc": {
        "_index": "_index",
        "_type": "_doc",
        "_id": "_id",
        "_source": {
          "my-keyword-field": "foo"
        },
        "_ingest": {
          "timestamp": "2099-03-07T11:04:03.000Z"
        }
      }
    },
    {
      "doc": {
        "_index": "_index",
        "_type": "_doc",
        "_id": "_id",
        "_source": {
          "my-keyword-field": "bar"
        },
        "_ingest": {
          "timestamp": "2099-03-07T11:04:04.000Z"
        }
      }
    }
  ]
}

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

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

POST my-data-stream/_doc?pipeline=my-pipeline
{
  "@timestamp": "2099-03-07T11:04:05.000Z",
  "my-keyword-field": "foo"
}

PUT my-data-stream/_bulk?pipeline=my-pipeline
{ "create":{ } }
{ "@timestamp": "2099-03-07T11:04:06.000Z", "my-keyword-field": "foo" }
{ "create":{ } }
{ "@timestamp": "2099-03-07T11:04:07.000Z", "my-keyword-field": "bar" }

Вы также можете использовать параметр pipeline с API обновления по запросу или API повторной индексации.

POST my-data-stream/_update_by_query?pipeline=my-pipeline

POST _reindex
{
  "source": {
    "index": "my-data-stream"
  },
  "dest": {
    "index": "my-new-data-stream",
    "op_type": "create",
    "pipeline": "my-pipeline"
  }
}

Установить конвейер по умолчанию

Используйте настройку индекса index.default_pipeline для установки конвейера по умолчанию. Elasticsearch применяет этот конвейер к запросам индексирования, если не указан параметр pipeline.

Установить конечный конвейер

Используйте настройку индекса index.final_pipeline для установки конечного конвейера. Elasticsearch применяет этот конвейер после запроса или конвейера по умолчанию, даже если ни один из них не указан.

Конвейеры для Beats

Чтобы добавить конвейер загрузки в Elastic Beat, укажите параметр pipeline в разделе output.elasticsearch в <BEAT_NAME>.yml. Например, для Filebeat вы бы указали pipeline в filebeat.yml.

output.elasticsearch:
  hosts: ["localhost:9200"]
  pipeline: my-pipeline

Конвейеры для Fleet и Elastic Agent

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

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

Fleet не предоставляет конвейер загрузки для интеграции Пользовательские логи. Вы можете безопасно указать конвейер для этой интеграции двумя способами: шаблон индекса или настройку конфигурации.

Вариант 1: Шаблон индекса

  1. Создайте и протестируйте ваш конвейер загрузки. Назовите свой конвейер logs-<dataset-name>-default. Это упростит отслеживание конвейера для вашей интеграции.

    Например, следующий запрос создаёт конвейер для набора данных my-app. Название конвейера — logs-my_app-default.

    PUT _ingest/pipeline/logs-my_app-default
    {
      "description": "Pipeline for `my_app` dataset",
      "processors": [ ... ]
    }
  2. Создайте шаблон индекса, который включает ваш конвейер в настройку индекса index.default_pipeline или index.final_pipeline. Убедитесь, что шаблон включён в поток данных. Шаблон индекса должен соответствовать logs-<dataset-name>-*.

    Вы можете создать этот шаблон с помощью функции Управление индексами Kibana или API создания шаблона индекса.

    Например, следующий запрос создает шаблон, соответствующий logs-my_app-*. Шаблон использует компонентный шаблон, который содержит настройку индекса index.default_pipeline.

    # Creates a component template for index settings
    PUT _component_template/logs-my_app-settings
    {
      "template": {
        "settings": {
          "index.default_pipeline": "logs-my_app-default",
          "index.lifecycle.name": "logs"
        }
      }
    }
    
    # Creates an index template matching `logs-my_app-*`
    PUT _index_template/logs-my_app-template
    {
      "index_patterns": ["logs-my_app-*"],
      "data_stream": { },
      "priority": 500,
      "composed_of": ["logs-my_app-settings", "logs-my_app-mappings"]
    }
  3. При добавлении или редактировании интеграции Пользовательские логи в Fleet, нажмите Настроить интеграцию > Файл пользовательского лога > Дополнительные параметры.
  4. В поле Имя набора данных укажите имя вашего набора данных. Fleet добавит новые данные для интеграции в полученный logs-<dataset-name>-default поток данных.

    Например, если имя вашего набора данных было my_app, Fleet добавляет новые данные в поток данных logs-my_app-default.

    Set up custom log integration in Fleet
  5. Используйте API изменения индекса для переключения вашего потока данных. Это гарантирует, что Elasticsearch применяет шаблон индекса и его настройки конвейера к любым новым данным для интеграции.

    POST logs-my_app-default/_rollover/

Вариант 2: Настройка по умолчанию

  1. Создайте и проверьте свой конвейер обработки данных. Назовите свой конвейер logs-<dataset-name>-default. Это упростит отслеживание конвейера для вашей интеграции.

    Например, следующий запрос создает конвейер для набора данных my-app. Название конвейера - logs-my_app-default.

    PUT _ingest/pipeline/logs-my_app-default
    {
      "description": "Pipeline for `my_app` dataset",
      "processors": [ ... ]
    }
  2. При добавлении или редактировании интеграции Пользовательские логи в Fleet, нажмите Настроить интеграцию > Файл пользовательского лога > Дополнительные параметры.
  3. В поле Название набора данных укажите имя вашего набора данных. Fleet добавит новые данные для интеграции в результирующий logs-<dataset-name>-default поток данных.

    Например, если название вашего набора данных было my_app, Fleet добавит новые данные в logs-my_app-default поток данных.

  4. В Пользовательские конфигурации укажите свой конвейер в параметре политики pipeline.

    Custom pipeline configuration for custom log integration

Установленный автономно Elastic Agent

Если вы используете Elastic Agent автономно, вы можете применять конвейеры, используя шаблон индекса, который включает настройку индекса index.default_pipeline или index.final_pipeline. В качестве альтернативы вы можете указать параметр политики pipeline в своей конфигурации elastic-agent.yml. См. Установка автономных Elastic Agent.

Доступ к исходным полям в процессоре

Процессоры имеют чтение и запись доступа к исходным полям входящего документа. Для доступа к ключу поля в процессоре используйте его имя поля. Следующий процессор set обращается к полю my-long-field.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "set": {
        "field": "my-long-field",
        "value": 10
      }
    }
  ]
}

Вы также можете добавить префикс _source.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "set": {
        "field": "_source.my-long-field",
        "value": 10
      }
    }
  ]
}

Используйте нотацию точки для доступа к полям объекта.

Если ваш документ содержит уплощенные объекты, используйте процессор dot_expander, чтобы расширить их сначала. Другие процессоры обработки данных не могут получить доступ к уплощенным объектам.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "dot_expander": {
        "description": "Expand 'my-object-field.my-property'",
        "field": "my-object-field.my-property"
      }
    },
    {
      "set": {
        "description": "Set 'my-object-field.my-property' to 10",
        "field": "my-object-field.my-property",
        "value": 10
      }
    }
  ]
}

Несколько параметров процессора поддерживают фрагменты шаблонов Mustache. Для доступа к значениям полей в фрагменте шаблона заключите имя поля в тройные фигурные скобки: {{{field-name}}}. Вы можете использовать фрагменты шаблонов для динамического задания имен полей.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "set": {
        "description": "Set dynamic '<service>' field to 'code' value",
        "field": "{{{service}}}",
        "value": "{{{code}}}"
      }
    }
  ]
}

Доступ к метаданным полей в процессоре

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

  • _index
  • _id
  • _routing
  • _dynamic_templates
PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "set": {
        "description": "Set '_routing' to 'geoip.country_iso_code' value",
        "field": "_routing",
        "value": "{{{geoip.country_iso_code}}}"
      }
    }
  ]
}

Используйте фрагмент шаблона Mustache для доступа к значениям метаданных. Например, {{{_routing}}} извлекает значение маршрутизации документа.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "set": {
        "description": "Use geo_point dynamic template for address field",
        "field": "_dynamic_templates",
        "value": {
          "address": "geo_point"
        }
      }
    }
  ]
}

Установленный выше процессор сообщает ES использовать динамический шаблон с именем geo_point для поля address, если это поле еще не определено в отображении индекса. Этот процессор переопределяет динамический шаблон для поля address, если он уже определен в запросе bulk, но не влияет на другие динамические шаблоны, определенные в запросе bulk.

Если вы автоматически генерируете идентификаторы документов, вы не можете использовать {{{_id}}} в процессоре. Elasticsearch назначает сгенерированные автоматически значения _id после обработки данных.

Доступ к метаданным обработки данных в процессоре

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

В отличие от исходных и метаданных полей, Elasticsearch по умолчанию не индексирует метаданные обработки данных. Elasticsearch также разрешает исходные поля, начинающиеся с ключа _ingest. Если ваши данные включают такие исходные поля, используйте _source._ingest для доступа к ним.

Конвейеры по умолчанию создают только метаданные обработки данных _ingest.timestamp. Это поле содержит отметку времени, когда Elasticsearch получил запрос индексирования документа. Для индексирования _ingest.timestamp или других полей метаданных обработки данных используйте процессор set.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "set": {
        "description": "Index the ingest timestamp as 'event.ingested'",
        "field": "event.ingested",
        "value": "{{{_ingest.timestamp}}}"
      }
    }
  ]
}

Обработка ошибок конвейера

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

Чтобы пропустить ошибку процессора и продолжить выполнение оставшихся процессоров конвейера, установите ignore_failure в значение true.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "rename": {
        "description": "Rename 'provider' to 'cloud.provider'",
        "field": "provider",
        "target_field": "cloud.provider",
        "ignore_failure": true
      }
    }
  ]
}

Используйте параметр on_failure, чтобы указать список процессоров, которые нужно выполнить сразу после ошибки процессора. Если on_failure указан, Elasticsearch после этого выполнит оставшиеся процессоры конвейера, даже если конфигурация on_failure пуста.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "rename": {
        "description": "Rename 'provider' to 'cloud.provider'",
        "field": "provider",
        "target_field": "cloud.provider",
        "on_failure": [
          {
            "set": {
              "description": "Set 'error.message'",
              "field": "error.message",
              "value": "Field 'provider' does not exist. Cannot rename to 'cloud.provider'",
              "override": false
            }
          }
        ]
      }
    }
  ]
}

Вложенный список процессоров on_failure для обработки ошибок вложенности.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "rename": {
        "description": "Rename 'provider' to 'cloud.provider'",
        "field": "provider",
        "target_field": "cloud.provider",
        "on_failure": [
          {
            "set": {
              "description": "Set 'error.message'",
              "field": "error.message",
              "value": "Field 'provider' does not exist. Cannot rename to 'cloud.provider'",
              "override": false,
              "on_failure": [
                {
                  "set": {
                    "description": "Set 'error.message.multi'",
                    "field": "error.message.multi",
                    "value": "Document encountered multiple ingest errors",
                    "override": true
                  }
                }
              ]
            }
          }
        ]
      }
    }
  ]
}

Вы также можете указать on_failure для конвейера. Если процессор без значения on_failure завершится ошибкой, Elasticsearch использует этот параметр на уровне конвейера как резервный вариант. Elasticsearch не будет пытаться выполнить оставшиеся процессоры конвейера.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [ ... ],
  "on_failure": [
    {
      "set": {
        "description": "Index document to 'failed-<index>'",
        "field": "_index",
        "value": "failed-{{{ _index }}}"
      }
    }
  ]
}

Дополнительная информация об ошибке конвейера может быть доступна в метаданных документа on_failure_message, on_failure_processor_type, on_failure_processor_tag и on_failure_pipeline. Эти поля доступны только внутри блока on_failure.

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

PUT _ingest/pipeline/my-pipeline
{
  "processors": [ ... ],
  "on_failure": [
    {
      "set": {
        "description": "Record error information",
        "field": "error_information",
        "value": "Processor {{ _ingest.on_failure_processor_type }} with tag {{ _ingest.on_failure_processor_tag }} in pipeline {{ _ingest.on_failure_pipeline }} failed with message {{ _ingest.on_failure_message }}"
      }
    }
  ]
}

Условное выполнение процессора

Каждый процессор поддерживает необязательное условие if, записанное как скрипт Painless. Если оно указано, процессор выполняется только тогда, когда условие if имеет значение true.

Скрипты условия if выполняются в контексте процессора обработки данных Painless Painless. В условиях if значения ctx являются только для чтения.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "drop": {
        "description": "Drop documents with 'network.name' of 'Guest'",
        "if": "ctx?.network?.name == 'Guest'"
      }
    }
  ]
}

Если включена настройка кластера script.painless.regex.enabled, вы можете использовать регулярные выражения в скриптах условия if. Для поддержки синтаксиса см. Регулярные выражения Painless.

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

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "set": {
        "description": "If 'url.scheme' is 'http', set 'url.insecure' to true",
        "if": "ctx.url?.scheme =~ /^http[^s]/",
        "field": "url.insecure",
        "value": true
      }
    }
  ]
}

Условие if должно быть указано как допустимый JSON на одной строке. Однако вы можете использовать синтаксис тройных кавычек консоли Kibana для написания и отладки более крупных скриптов.

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

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "drop": {
        "description": "Drop documents that don't contain 'prod' tag",
        "if": """
            Collection tags = ctx.tags;
            if(tags != null){
              for (String tag : tags) {
                if (tag.toLowerCase().contains('prod')) {
                  return false;
                }
              }
            }
            return true;
        """
      }
    }
  ]
}

Вы также можете указать сохраненный скрипт как условие if.

PUT _scripts/my-prod-tag-script
{
  "script": {
    "lang": "painless",
    "source": """
      Collection tags = ctx.tags;
      if(tags != null){
        for (String tag : tags) {
          if (tag.toLowerCase().contains('prod')) {
            return false;
          }
        }
      }
      return true;
    """
  }
}

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "drop": {
        "description": "Drop documents that don't contain 'prod' tag",
        "if": { "id": "my-prod-tag-script" }
      }
    }
  ]
}

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

Например, ctx.network?.name.equalsIgnoreCase('Guest') не является безопасным для NULL. ctx.network?.name может возвращать NULL. Перепишите скрипт как 'Guest'.equalsIgnoreCase(ctx.network?.name), который является безопасным для NULL, потому что Guest всегда не NULL.

Если вы не можете переписать скрипт для работы со случаями NULL, добавьте явную проверку на NULL.

PUT _ingest/pipeline/my-pipeline
{
  "processors": [
    {
      "drop": {
        "description": "Drop documents that contain 'network.name' of 'Guest'",
        "if": "ctx.network?.name != null && ctx.network.name.contains('Guest')"
      }
    }
  ]
}

Условное применение конвейеров

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

PUT _ingest/pipeline/one-pipeline-to-rule-them-all
{
  "processors": [
    {
      "pipeline": {
        "description": "If 'service.name' is 'apache_httpd', use 'httpd_pipeline'",
        "if": "ctx.service?.name == 'apache_httpd'",
        "name": "httpd_pipeline"
      }
    },
    {
      "pipeline": {
        "description": "If 'service.name' is 'syslog', use 'syslog_pipeline'",
        "if": "ctx.service?.name == 'syslog'",
        "name": "syslog_pipeline"
      }
    },
    {
      "fail": {
        "description": "If 'service.name' is not 'apache_httpd' or 'syslog', return a failure message",
        "if": "ctx.service?.name != 'apache_httpd' && ctx.service?.name != 'syslog'",
        "message": "This pipeline requires service.name to be either `syslog` or `apache_httpd`"
      }
    }
  ]
}

Получение статистики использования конвейера

Используйте API статистики узлов, чтобы получить глобальную и по-конвейерную статистику по загрузке. Используйте эту статистику, чтобы определить, какие конвейеры выполняются чаще всего или тратят больше всего времени на обработку.

GET _nodes/stats/ingest?filter_path=nodes.*.ingest

© 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/7.17/ingest.html

Spec-Zone.ru

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