Заявка на потоковую обработку#

Заявка на потоковую обработку создается пользователем — представителем системы/сервиса-получателя данных для передачи преобразованных ранее опубликованных событий из топика домена-источника в топик домена-получателя.

Вводная информация о потоковой обработке#

Технический сервис потоковой обработки служит для выполнения преобразований ранее опубликованных событий. Представляет собой Apache Flink job, который может подписываться на потоки событий, выполнять преобразования и в результате публиковать новые потоки событий.

Обработчик устанавливается в рамках домена, в одном домене может быть только один обработчик. Топики-источники, обработчик и топики-назначения могут быть в разных доменах. После установки обработчика необходимо добавить его в справочник Потоковые обработчики в EDMS (выполняет оператор с ролью EMC Admin). После этого возможно создание заявок на потоковую обработку пользователем с ролью System User, System Owner и правами на использование системы.

Пример схемы размещения компонентов при потоковой обработке событий
#

Replication Scheme

Создание новых конфигураций обработчиков#

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

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

Более подробно формирование conf и tr файлов описано в документации к компоненту EVTP в секциях документа Руководство пользователя EVTP разделы: Построение потока обработки и DSL функции. Ниже приведена общая информация по формированию conf и tr файлов.

Формирование файла *.conf#

Файл *.conf обязательно должен содержать следующие секции:

  • наименование топика и перечень источников;

  • наименование топика и перечень источников:

  flow: {
    name: "<наименование>"
    source: [${<наименование источника 1>},${<наименование источника n>}]
  }
  • определение источника;

  <наименование источника>: {
    name: "<описание источника в формате Сегмент.Домен.Имя_потока_в_правильном_регистре"
    type: "source"
    topic: "<наименование topic>"
    config: ${<ссылка на конфигурацию транспорта из defaults.json>} {
      "consumer": {
        "group.id": "<consumer group, с которой производится подключение>"
      }
    }
    destination: ${<наименование следующего шага обработки>}
  }
  • определение выходной точки;

  <наименование выходной точки>: {
    name: "<описание выходной точки Сегмент.Домен.Имя_потока_в_правильном_регистре>"
    type: "destination"
    topic: "<наименование topic>"
    config: ${<ссылка на конфигурацию транспорта из defaults.json>}
  }
  • определение шагов преобразований (дополнительная секция).

  <наименование шага преобразования>: {
    name: "<описание шага преобразования>"
    type: "dsl"
    format: {
      input: "<указание входного формата xml/json>"
      output: "<указание выходного формата xml/json>"
    }

    dsl: {
      source: "file"
      path: "путь до tr файла, описывающего преобразования"
    }
    destination: ${<наименование следующего шага обработки>}
  }

В зависимости от типа развертывания, используются следующие варианты задания пути до tr файла:

  • для прямого развертывания: "config/<наименование файла>/<наименование tr файла>" (Пример: "config/test/test_step.tr")

  • для развертывания с использованием конфигурационных дистрибутивов: ${defaults.flinkInstalldir}"/mapping/<наименование файла>/<наименование tr файла>" (Пример: ${defaults.flinkInstalldir}"/mapping/test/test_step.tr")

Формирование файла *.tr#

В файле *.tr содержатся необходимые преобразования событий. По умолчанию обработчик не производит никаких действий над потоком, все необходимые действия должны указываться явно.

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

  • копирование заголовков:

  headers.kafka = Environment.kafka.headers
  • копирование тела сообщения:

  OUT = IN
  • передача на выход тела сообщения, полученного на вход:

  return("passthrough")
  • журналирование:

  INFO(<текст>, <идентификатор события>, <источник события>, <получатель события>)

Порядок действий для обеспечения процесса трансформации события#

  1. Развернуть в домене назначения топик/топики, в который(е) обработчик будет записывать обработанное событие.

Выполняется по заявке на публикацию. Заявка(и) создается автоматически вместе с заявкой на потоковую обработку. Количество заявок на публикацию соответствует количеству исходящих потоков событий.

  1. Выполнить подписку обработчика на топик-источник в домене-источнике.

Выполняется по заявке на подписку. Заявка(и) создается автоматически вместе с заявкой на потоковую обработку. Количество заявок на подписку соответствует количеству входящих потоков событий.

  1. Запустить потоковый обработчик, который после преобразования будет публиковать событие из топика(ов)-источник(ов) в топик(и)-назначения.

Выполняется по заявке на потоковую обработку, созданной пользователем с ролью System User, System Owner.

Создание и заполнение заявки#

  1. В разделе Заявки перейти на вкладку Заявки на потоковую обработку, нажать кнопку Создать заявку на потоковую обработку:

Create Transformation Request

Вкладка Обработчик.#

Здесь содержится информация про обработчик для данной интеграции.

В данной вкладке необходимо заполнить информацию о потоковом обработчике для данной интеграции.

Transformation Parameters

  • Система-инициатор. Выбрать из списка код системы, которую представляет автор заявки. В поле реализован поиск по коду системы.

Если у вас есть доступ к этой системе, будет выполнен поиск по введенным символам. Если система не отображается в выпадающем списке: Пользователю с ролью System User, необходимо обратиться к System Owner для добавления в список пользователей системы. Пользователю с ролью System Owner, необходимо обратиться к EMC Admin для добавления в список пользователей системы.

  • Сегмент. Сегмент, в котором происходит потоковая обработка.

  • Домен. Содержит код домена, в котором располагается обработчик событий. В случае, если обработчик не установлен (не зарегистрирован) в данном домене, то будет отображаться информационная строка:

Transformation Parameters

  • Имя обработчика. Имя обработчика, заполняется пользователем с ролью System User, System Owner. Правило наименования: <СЕГМЕНТ>.<КОД ДОМЕНА>.<НАИМЕНОВАНИЕ conf ФАЙЛА>. Имя обработчика должно соответствовать значению name в блоке flow» конфигурационного файла обработчика.

  • Конфигурация обработчика. При нажатии на кнопку необходимо загрузить конфигурацию обработчика в формате zip. В архиве должны содержаться несколько обязательных файлов - один с расширением .conf и произвольное количество с расширением .tr, в зависимости от количества шагов обработки. Подробное описание структуры конфигурации представлено в документации к компоненту EVTP в секциях документа Руководство пользователя EVTP разделы: Построение потока обработки и DSL функции. Для замены ранее загруженной схемы необходимо снова нажать Конфигурация обработчика и загрузить новую конфигурацию - доступно только для заявки в статусе Черновик, которой еще не присвоен контур. Для скачивания схемы необходимо нажать на кнопку Transformation Parameters.

  • Описание потоковой обработки. Необходимо добавить краткое описание обработки

  • Бизнес-процесс. В данном поле необходимо указать идентификатор процесса, в рамках которого создана текущая заявка. Код интеграции можно уточнить у архитектора интеграции.

Вкладка События для потоковой обработки.
#

В данной вкладке содержится информация о событиях, которые будут обрабатываться.

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

Create new system

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

Transformation Parameters

Названия выбранных событий-источников добавить в поле topic раздела source в конфигурационном файле обработчика, который будет загружаться на вкладке Обработчик данной заявки. Подробное описание представлено в документации к компоненту EVTP в секциях документа Руководство пользователя EVTP разделы: Построение потока обработки

  • Группа-потребителей. Заполняется пользователем с ролью System User, System Owner. В зависимости от настроек EDMS, группа может быть либо самостоятельно введена пользователем, либо система предложит шаблонную группу для выбора. Доступно для редактирования. Важно, чтобы группа потребителей соответствовала описанию в .conf файле, но может отличаться для разных событий.

Названия групп необходимо добавить в поле group.id раздела source/consumer в конфигурационном файле обработчика, который будет загружаться на вкладке Обработчик данной заявки. Подробное описание представлено в документации к компоненту EVTP в секциях документа Руководство пользователя EVTP разделы: Построение потока обработки

Transformation Parameters

Просмотр информации об источнике событий:

Transformation Parameters

Вкладка Результирующие события.#

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

В текущей вкладке заполняются только поля результирующих событий. Поля: домен и имя топика сохраняются из предыдущих вкладок.

Transformation Parameters

  • Название события. Название результирующего события, заполняется пользователем с ролью System User, System Owner. Название поля должно соответствовать значению topic из блока destination конфигурационного файла обработчика. Подробное описание представлено в документации к компоненту EVTP в секциях документа Руководство пользователя EVTP разделы: Построение потока обработки

  • Имя топика. Формируется автоматически по правилу CEP.<Название события>EVENT. Недоступно для редактирования. Имя топика является уникальным в рамках домена, таким образом, если другой пользователь уже создавал заявку на публикацию события с такими же именем топика — повторно это сделать невозможно.

  • Домен. Код домена выбирается из списка. В поле Домен реализован поиск по коду домена. При заполнении данного поля над вкладками появится информационное сообщение об ограничениях на домене. Эта информация доступна, если данные были заполнены пользователем с ролью Domain Admin. Если данный домен имеет статус Выводится из эксплуатации или Архив, при попытке передачи заявки на согласование появится ошибка: «Невозможна работа с заявкой: домен не активен». Для активации домена обратитесь к пользователю с ролью Domain Admin или EMS Admin.

  • Описание события. Текстовое поле с кратким описанием бизнес-смысла события. Ограничение на количество символов — 1024.

  • Формат схемы. Выбирается из списка. Возможные значения: json, avro, xml. Выберите схему. При нажатии на кнопку Схема события можно загрузить схему формата события. Для замены ранее загруженной схемы необходимо снова нажать Схема события и загрузить новую схему - доступно только для заявки в статусе Черновик, которой еще не присвоен контур. Для скачивания схемы необходимо нажать на кнопку Transformation Parameters.

  • Время хранения данных в топике (минуты) — дни, часы, минуты. Для снятия ограничения нажмите чекбокс Время хранения не ограничено.

  • Максимальный размер сообщения (Килобайты). Ограничение — 5000 Кб.

  • Пиковое значение TPS (запрос/секунда) — пиковое количество запросов в секунду.

  • Пиковое значение TPD (запрос/сутки) — ожидаемое пиковое количество запросов в сутки.

  • Количество партиций поле заполняется автоматически. Доступно для редактирования. Рекомендуется использовать значение для количества партиций кратное 2.

  • Фактор репликации — по умолчанию 2. Ограничения от 1 до 6. Указывается количество копий каждой секции топика, которые хранятся в кластере Kafka. Каждая секция топика реплицируется на указанное количество брокеров.

  • Максимальный размер данных на партицию в топике (Килобайты). Ограничение — 51200. Доступно для редактирования. Для снятия ограничения нажмите чек бокс Размер данных не ограничен.

  • Чекбокс Время хранения не ограничено позволяет установить неограниченное время хранения данных в топике событий. При включении появится информационное сообщение: Неограниченное время хранения приведет к увеличению КТС.

  • Чекбокс Размер данных не ограничен позволяет установить неограниченный размер данных в топике. По умолчанию выключен. При включении появится информационное сообщение: Неограниченный размер приведет к увеличению КТС.

  • Чекбокс Свертывание по ключу события (compaction) — данный переключатель используется в случае публикации в компактный топик событий. Compaction - процесс очистки лога в Kafka, сохраняющий последнее значение для каждого ключа в топике. Используется для экономии места и восстановления состояния. По умолчанию выключен.

Добавить новое результирующее событие можно по кнопке Добавить событие, удалить — по кнопке Удалить событие.

После заполнения заявки, сохраните ее, нажав кнопку Сохранить. Заявка сохраняется в статусе Черновик. Заявку в статусе Черновик можно отредактировать или удалить (кнопка Delete в каталоге заявок). После сохранения заявки будет присвоен идентификационный номер, сквозной в текущей инсталляции EDMS.

После сохранения заявки на потоковую обработку автоматически будут созданы связанные с ней заявки:

  • на подписку обработчика.

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

Transformation Parameters

Созданные автозаявки находятся в соответствующих разделах EDMS: Заявки на подписку и Заявки на публикацию. Также в заявке на потоковую обработку есть ссылки на созданные автозаявки, в блоке Связанные заявки, расположенном сверху внутри заявки. Из главной заявки можно перейти к любой из автоматически созданных заявок, развернув блок Связанные заявки и нажав на название заявки. Из автоматически связанной заявки можно вернуться к главной заявке. Кнопка К заявке на потоковую обработку.

Нажать кнопку Передать на согласование и отправить ее на согласование владельцу и администратору домена назначения.

Кнопка Архивировать заявку. Подробнее можно ознакомиться в разделе Роль - Пользователь с правами редактирования (System User).

Автоматическая заявка на публикацию потока назначения#

Автоматическая заявка на публикацию потока назначения размещается в разделе Заявка на публикацию, как и обычные заявки на публикацию. Таких заявок может быть несколько.

Она доступна для просмотра, но недоступна для редактирования вручную: вкладки *Свойства события, Схема и История настроек транспорта - все данные в ней заполняются автоматически на основании данных из заявки на на потоковую обработку. Вкладка Настройка транспорта частична доступна для редактирования вручную. В ней можно изменять следующие поля: Время хранения данных в топике (минуты), Максимальный размер данных на партицию в топике (Килобайты) и Количество партиций.

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

Вкладка Свойства события:

Заполнена данными о событии, которое будет публиковаться в поток назначения в результате потоковой обработки.

Вкладка Схема:

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

Вкладка Настройка транспорта:

Содержит параметры нагрузки, аналогичные параметрам, заданным для топика-источника.

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

При необходимости можно изменить следующие поля: Время хранения данных в топике (минуты), Максимальный размер данных на партицию в топике (Килобайты) и Количество партиций.

Вкладка История настроек транспорта

На вкладке отображаются контуры, на которые установлена заявка, а также сертификаты и параметры настройки транспорта на данном контуре.

Эти данные становятся доступны после перехода заявки в статус «Активно — установка выполнена» или «Активно».

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

Автоматическая заявка на подписку обработчика на поток-источник#

Автоматические заявки на подписку обработчика на поток-источник размещаются в разделе Заявка подписку, как и обычные заявки на подписку. Таких заявок может быть несколько.

Они доступны для просмотра, но недоступны для редактирования вручную - все данные в них заполняются автоматически на основании данных из заявки на потоковую обработку.

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

Вкладка Описание события:

Заполнена данными о событии, которое было выбрано для потоковой обработки.

Вкладка Подписчик

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

Вкладка Настройка транспорта

Значение Группа потребителей сформировано автоматически на основании данных о топике-источнике и топике назначения из заявки на потоковую обработку.

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

Управление заявкой после создания#

После создания и сохранения заявки вы можете выполнить следующие действия: