[Документация Yandex Cloud](../../index.md) > [Yandex Serverless Integrations](../index.md) > Практические руководства > Миграция с EventRouter на триггеры

# Миграция с EventRouter на триггеры

В качестве альтернативы EventRouter вы можете использовать:
* [триггеры](../../functions/concepts/trigger/index.md) для вызова функций Cloud Functions;
* [триггеры](../../serverless-containers/concepts/trigger/index.md) для вызова контейнеров Serverless Containers;
* [триггеры](../../api-gateway/concepts/trigger/index.md) для отправки событий в WebSocket-соединения;
* триггеры для запуска рабочих процессов Workflows.

В отличие от EventRouter, где шина — это набор правил и коннекторов, триггер — это один источник и один или несколько приемников.

## Что переносится автоматически {#auto-migration}

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

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

* в правилах шины заданы фильтр или шаблон преобразования;
* коннектор имеет тип **Audit Trails** или **API EventRouter**;
* среди приемников любого правила шины есть поток данных Yandex Data Streams, лог-группа Yandex Cloud Logging или очередь сообщений Yandex Message Queue;
* статус коннектора — не `Запущен` и не `Остановлен`;
* у коннектора с типом **Таймер** задан часовой пояс.

{% note warning %}

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

{% endnote %}

## Соответствие сущностей EventRouter и триггеров {#entity-mapping}

EventRouter | Триггеры | Что учесть при переносе
 --- | --- | --- 
Шина | Нет аналога | Триггер связывает источник и приемники напрямую, промежуточной шины нет.
Коннектор | Источник триггера | Один коннектор соответствует одному триггеру.
Правило | Приемник триггера | Правила принадлежат шине, а не коннектору: каждое событие проходит через все правила шины. Поэтому в каждый триггер попадают приемники всех правил шины.
Фильтр правила | Фильтр приемника | В EventRouter фильтр один на все приемники правила, в триггере фильтр задается отдельно для каждого приемника.
Приемник | Приемник триггера | Не более 5 приемников для одного триггера. Считаются суммарно по всем правилам шины.
Шаблон преобразования приемника | Шаблон преобразования приемника | Применяется к событию другого формата, подробнее в [Форматы сообщений](#message-format).
Настройки группирования | Настройки группирования для источника | В EventRouter настройки группирования задаются для каждого приемника, в триггере — одни на все приемники. Максимальный размер группы уменьшается с 256 КиБ до 64 КиБ.
Число повторных попыток отправки | Число повторных попыток отправки | В EventRouter — от 0 до 10, в триггерах — от 1 до 5. Для триггера с источником Yandex Message Queue повторные попытки отправки недоступны.
Нет аналога | Интервал между повторными попытками | В EventRouter не настраивается, в триггерах задается в диапазоне от 10 секунд до 1 минуты.
Максимальный срок жизни события | Нет аналога | В EventRouter событие перенаправляется в Dead Letter Queue, когда его возраст превышает заданное значение. В триггерах эта настройка отсутствует.
Dead Letter Queue приемника | Dead Letter Queue приемника | Для триггера с источником Yandex Message Queue недоступна: вместо нее используйте политику перенаправления самой очереди.
Часовой пояс таймера | Нет аналога | Расписание триггера задается только по UTC\+0.
Настройки чтения очереди | Частично | Таймаут видимости переносится, у размера группы при чтении и таймаута опроса аналогов нет.
Защита от удаления | Нет аналога | —
Логирование шины | Нет аналога | —
Остановка коннектора, отключение правила | Приостановка триггера | Для очереди сообщений и потока данных события копятся и обрабатываются после возобновления работы триггера. Для таймера и других источников без буфера события за время простоя теряются.
API EventRouter | Нет аналога | Подробнее читайте в [API EventRouter и прямая отправка в шину](#api-connector).

## План миграции {#migration-plan}

### Шаг 1. Составьте список ресурсов {#step-1-list-resources}

Составьте список всех шин, коннекторов и правил, которые нужно перенести:

```bash
yc serverless eventrouter bus list
yc serverless eventrouter connector list
yc serverless eventrouter rule list
```

Команды `connector list` и `rule list` выводят все коннекторы и правила каталога — отберите нужные по полю `bus_id` в выводе. Фильтровать по шине с помощью команды нельзя.

Для каждой шины зафиксируйте ее коннекторы с настройками источников и все ее правила с фильтрами и приемниками. Каждый коннектор станет отдельным триггером, а приемники **всех правил шины** — приемниками каждого из этих триггеров.

{% note warning %}

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

{% endnote %}

### Шаг 2. Проверьте ограничения {#step-2-check-limits}

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

* источник — API EventRouter или события отправляются в шину напрямую: нужно доработать приложение-отправитель;

* среди приемников есть очередь сообщений Message Queue, поток данных Data Streams или лог-группа Cloud Logging: нужна функция-прослойка;

* в правилах шины суммарно больше 5 приемников. Что делать, зависит от источника:

  * поток данных Data Streams — заведите в потоке данных дополнительного потребителя и создайте второй триггер с тем же потоком данных. Каждый потребитель получает полную копию событий, так что приемники можно разложить по нескольким триггерам. Указывать в двух триггерах одного потребителя нельзя: тогда они поделят события между собой;
  * таймер — создайте второй триггер с тем же расписанием;
  * очередь сообщений Message Queue — несколько триггеров на одну очередь создать нельзя, поэтому лишние вызовы придется вынести в функцию-прослойку, которая вызовет остальные приемники;

* в таймере указан часовой пояс или секунды в cron-выражении;

* у приемников заданы разные настройки группирования: в триггере они общие для всех приемников, поэтому придется выбрать одни. Автоматический перенос в таких случаях берет минимальные значения по всем приемникам;

* суммарный размер группы превышает 64 КиБ: в EventRouter лимит 256 КиБ;

* у приемников настроены повторные вызовы или Dead Letter Queue, а источник — очередь сообщений: для такого триггера ни то, ни другое не поддерживается, повторы настраиваются политикой перенаправления самой очереди, и для этого нужна роль `ymq.admin`;

* размер отдельного события превышает 230 КБ.

{% note warning %}

Событие больше допустимого размера триггер отбрасывает без повторной попытки и без записи в Dead Letter Queue. Убедитесь, что таких событий в источнике нет.

{% endnote %}

Убедитесь, что в облаке не более 100 триггеров. Эта квота общая для триггеров Cloud Functions, Serverless Containers и API Gateway, и шина с несколькими коннекторами расходует ее быстро. Чтобы повысить квоту, обратитесь в [техническую поддержку](https://kz.center.yandex.cloud/support).

### Шаг 3. Подготовьте триггеры к переключению {#step-3-prepare-triggers}

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

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

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

1. Убедитесь, что формат событий соответствует ожиданиям, и подготовьте jq-шаблоны и изменения в коде приемников.

### Шаг 4. Переключитесь на триггеры {#step-4-switch}

Переключение зависит от типа источника:

Источник | Порядок действий | Что происходит с событиями
--- | ---| --- |
Yandex Message Queue | Остановите коннектор, дождитесь обработки накопленных событий, создайте триггер | События сохраняются в очереди. Коннектор и триггер могут читать одну очередь одновременно, но тогда они поделят события между собой, поэтому проверить два контура параллельно на одной очереди нельзя — для проверки нужна вторая очередь.
Yandex Data Streams | Остановите коннектор, создайте триггер с тем же потребителем | События сохраняются в потоке. Если создать триггер с новым потребителем, чтение начнется не с того же места.
Таймер | Остановите коннектор, создайте триггер | Одно срабатывание может быть пропущено или продублировано.
API EventRouter, прямая отправка в шину | Отправьте из приложения событие в очередь сообщений или поток данных, создайте триггер | События, отправленные в шину после остановки, теряются.

Триггеры, как и EventRouter, гарантируют доставку `At least once`, поэтому во время переключения возможны повторные вызовы. Убедитесь, что обработчики идемпотентны.

### Шаг 5. Проверьте работу {#step-5-verify}

* Отправьте тестовое событие и убедитесь, что приемник вызван.
* Сравните метрики вызовов приемников до и после переключения.
* Если в настройках приемника указана Dead Letter Queue, проверьте, что она не пополняется. Для триггера с источником Yandex Message Queue такой проверки нет — смотрите DLQ, указанную в настройках политики перенаправления очереди.

### Шаг 6. Удалите ресурсы EventRouter {#step-6-cleanup}

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

## Пример миграции шины {#migration-example}

Ниже разобран типовой случай: шина с одним коннектором и двумя правилами превращается в один триггер с двумя приемниками.

### Что было в EventRouter {#example-before}

К шине `orders-bus` подключен один коннектор `orders-queue` с источником Yandex Message Queue: очередь `orders`. События в очереди выглядят так:

```json
{"orderId": "1234", "status": "new", "amount": 500}
```

К шине привязаны два правила, у обоих приемников задано группирование по 10 событий или 5 секунд:

Правило | Фильтр | Приемник
--- | --- | ---
`process-orders` | `.status == "new"` | Контейнер `order-processor`
`notify-orders` | Не задан | Функция `order-notifier`

### Что получится в триггерах {#example-after}

Один триггер с источником Yandex Message Queue и двумя приемниками:

Было | Стало
--- | ---
Коннектор `orders-queue` | Источник триггера: очередь `orders`
Настройки группирования у каждого приемника | Одни настройки группирования на источнике: 10 сообщений или 5 секунд
Правило `process-orders` | Приемник 1: вызов контейнера `order-processor` с фильтром
Правило `notify-orders` | Приемник 2: вызов функции `order-notifier`

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

Если бы приемником одного из правил была лог-группа, очередь сообщений или поток данных, вместо приемника понадобилась бы [функция-прослойка](#shim) — прямых аналогов у этих приемников нет.

### Шаг 1. Опишите действия {#example-step-1}

Фильтр правила `process-orders` нельзя перенести дословно: он был написан для тела события, а приемник триггера получает JSON-объект с событием. По [рецепту переноса](#message-format-jq) допишите слева распаковку тела. Тем же выражением задайте шаблон, чтобы контейнер получал внутри JSON-объекта тела событий, а не служебную обертку сообщения очереди.

Приемник 1 — вызов контейнера.

```json
{
  "invokeContainer": {
    "containerId": "<идентификатор_контейнера_order-processor>",
    "serviceAccountId": "<идентификатор_сервисного_аккаунта>"
  },
  "filter": {"jq": ".details.message.body | fromjson | .status == \"new\""},
  "transformer": {"jq": ".details.message.body | fromjson"}
}
```

Приемник 2 — вызов функции. У правила `notify-orders` фильтра не было, поэтому и в приемнике его нет.

```json
{
  "invokeFunction": {
    "functionId": "<идентификатор_функции_order-notifier>",
    "serviceAccountId": "<идентификатор_сервисного_аккаунта>"
  },
  "transformer": {"jq": ".details.message.body | fromjson"}
}
```

Ни в одном из приемников нет ни повторных вызовов, ни Dead Letter Queue: источник — очередь сообщений, а для такого триггера они не поддерживаются. Если указать для приемника `retryPolicy` или `deadLetter`, создание триггера завершится ошибкой. Повторную обработку настраивайте с помощью политики перенаправления самой очереди — для этого нужна роль `ymq.admin`.

### Шаг 2. Создайте триггер {#example-step-2}

Сохраните описания приемников в файлы `action-1.json` и `action-2.json` и создайте триггер:

```bash
yc serverless trigger v2 create message-queue orders \
  --queue-arn <ARN_очереди> \
  --service-account-id <идентификатор_сервисного_аккаунта> \
  --batch-max-count 10 \
  --batch-cutoff 5s \
  --action @action-1.json \
  --action @action-2.json
```

Параметр `--action` можно передать строкой или ссылкой на файл через `@`. Рекомендуем второй способ: jq-выражения содержат кавычки, и во встроенном JSON их приходится экранировать.

Готовый шаблон можно получить с помощью команды `yc serverless trigger v2 help-action --invoke-container`. Аналогично для `--invoke-function`, `--start-workflow` и `--gateway-websocket-broadcast`.

{% note warning %}

Обязательно указывайте `v2` в пути команды. Без него вызывается устаревшая группа команд `yc serverless trigger v1`, которая пока остается вариантом по умолчанию и не поддерживает `--action`, фильтры, шаблоны и выбор потребителя. Попытка выполнить команду без `v2` завершится ошибкой `unknown flag: --action`.

{% endnote %}

Триггер с несколькими приемниками, фильтрами и шаблонами нельзя создать отдельными параметрами вида `--invoke-function-id` — они задают один приемник без дополнительных настроек. Используйте консоль управления, `yc serverless trigger v2`, API v2 или Terraform.

На каждый приемник приходится один параметр `--action`, в триггере может быть не больше пяти приемников.

### Что изменится для приемников {#example-consumer-changes}

Группирование было включено и раньше, поэтому контейнер `order-processor` уже получал не отдельное событие, а JSON-массив тел:

```json
[
  {"orderId": "1234", "status": "new", "amount": 500}
]
```

Он продолжит получать только события со статусом `new`, но теперь массив будет лежать в JSON-объекте по ключу `messages`:

```json
{
  "messages": [
    {"orderId": "1234", "status": "new", "amount": 500}
  ]
}
```

Код контейнера нужно научить разворачивать JSON-объект: [убрать его шаблоном нельзя](#message-format).

То же касается функции `order-notifier`: она получит все события пакетом, в JSON-объекте, и по-прежнему без фильтрации.

### Если источник — поток данных {#example-yds-source}

Для коннектора с источником Yandex Data Streams порядок тот же, с двумя отличиями:

* при создании триггера укажите того же потребителя, который был задан в коннекторе, иначе чтение начнется не с того же места;
* фильтр и шаблон переносятся без изменений — элементы JSON-объекта совпадают с записями потока, распаковывать тело не нужно. Фильтр правила остается выражением `.status == "new"`, а шаблон не нужен.

## Миграция источников (коннекторов) {#source-migration}

### Таймер {#source-timer}

Создайте [таймер](../../functions/concepts/trigger/timer.md).

Cron-выражение из коннектора нельзя перенести в триггер без изменений — в EventRouter и в триггерах разный порядок полей. Скрытой подмены расписания при этом не произойдет: триггер откажется принять скопированное выражение. В EventRouter одно из полей `Day of month` и `Day of week` всегда содержит `?`, а при сдвиге полей этот символ попадает в `Month` или `Year`, где он недопустим. Создание триггера завершится ошибкой вида `'?' can only be specified for Day-of-Month or Day-of-Week`. Преобразуйте выражение по таблице ниже.

Функциональность | Порядок полей в cron-выражении
--- | ---
EventRouter | `Seconds Minutes Hours Day-of-month Month Day-of-week [Year]`
Триггеры | `Minutes Hours Day-of-month Month Day-of-week [Year]`

Чтобы преобразовать выражение, уберите первое поле `Seconds`. Поле `Year` необязательно в обеих функциональностях: если оно было задано в коннекторе, перенесите его без изменений, а если нет — можно оставить выражение из пяти полей или дописать `*`.

Примеры cron-выражений:

EventRouter | Триггеры | Описание
--- | --- | ---
`0 * * * * ?` | `* * * * ? *` | Каждую минуту
`0 0 * ? * *` | `0 * ? * * *` | Каждый час
`0 15 10 ? * *` | `15 10 ? * * *` | Каждый день в 10:15

Как и в EventRouter, поля `Day of month` и `Day of week` нельзя заполнять одновременно: если значение задано в одном, во втором должен стоять `?`. При переносе следите, чтобы `?` не потерялся вместе со сдвигом полей.

Нумерация дней недели в обоих сервисах одинаковая — `1` соответствует воскресенью, `7` — субботе.

{% note warning %}

Триггеры не поддерживают указание секунд в cron-выражении. Минимальная единица измерения — 1 минута. Если сценарий требует более частого срабатывания, пересмотрите логику работы приложения.

{% endnote %}

{% note warning %}

В триггерах нельзя задать часовой пояс, время в cron-выражении всегда указывается по UTC\+0. Если в коннекторе был задан другой часовой пояс, пересчитайте время в расписании самостоятельно. Учтите, что при таком пересчете расписание перестанет автоматически учитывать переход на летнее и зимнее время, если он есть в вашем часовом поясе.

{% endnote %}

### Yandex Message Queue {#source-ymq}

Создайте [триггер для Yandex Message Queue](../../functions/concepts/trigger/ymq-trigger.md).

Формат сообщения от триггера для Message Queue отличается от формата в EventRouter. Подробнее в [Форматы сообщений](#message-format). Рекомендуем использовать шаблон преобразования в настройках триггера или изменить конфигурацию вызываемых ресурсов для адаптации под ваши задачи. Например, для получения только тела сообщения укажите в настройках триггера шаблон `.details.message.body`.

{% note warning %}

У приемников такого триггера не поддерживаются повторные попытки отправки и Dead Letter Queue. Повторную обработку настраивайте с помощью политики перенаправления самой очереди — для этого нужна роль `ymq.admin`.

{% endnote %}

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

Обратите внимание, что в качестве источника триггера задается ARN очереди — так же, как в приемнике EventRouter. URL очереди понадобится только [функции-прослойке](#shim), если очередь была еще и приемником.

### Yandex Data Streams {#source-yds}

Создайте [триггер для Yandex Data Streams](../../functions/concepts/trigger/data-streams-trigger.md).

Чтобы избежать повторной обработки событий, убедитесь, что при создании триггера указан **тот же потребитель**, который был настроен в коннекторе. В устаревшей группе команд `yc serverless trigger` такого поля нет: сервис заведет собственного потребителя с именем по идентификатору триггера, и позиция чтения потеряется. Указывайте потребителя через `yc serverless trigger v2 create yds` или API v2.

Содержимое записей триггер передает без изменений, но оборачивает их в JSON-объект `{"messages": [...]}`. Подробнее в [Форматы сообщений](#message-format).

### Audit Trails {#source-audit-trails}

Прямого триггера для событий Audit Trails в текущей реализации не предусмотрено. Если вы используете EventRouter для обработки аудитных событий сервисов Container Registry или Object Storage, подойдут [триггер для Container Registry](../../functions/concepts/trigger/cr-trigger.md) или [триггер для Object Storage](../../functions/concepts/trigger/os-trigger.md).

Для других сценариев необходимо:

1. Экспортировать события в [поток данных](../../data-streams/concepts/glossary.md#stream-concepts).
1. Настроить интеграцию с Audit Trails, [создав трейл](../../audit-trails/concepts/trail.md) и указав созданный поток данных как объект назначения.
1. [Создать триггер для Yandex Data Streams](../../functions/concepts/trigger/data-streams-trigger.md), указав созданный поток данных как источник.

### API EventRouter и прямая отправка в шину {#api-connector}

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

* через коннектор с типом источника **API EventRouter** — вызов `EventService/Send` или команда `yc serverless eventrouter send-event`;
* напрямую в шину, не используя коннектор, — вызов `EventService/Put` или команда `yc serverless eventrouter put-event`.

Оба способа перестанут работать. Чтобы триггер запускался по событиям, которые генерирует ваше приложение, приложение должно записывать их напрямую в очередь Yandex Message Queue или поток данных Yandex Data Streams, а триггер — читать из этой очереди или потока.

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

#### Вариант 1: Отправка через Yandex Data Streams {#api-option-yds}

1. Создайте [поток данных Yandex Data Streams](../../data-streams/concepts/glossary.md#stream-concepts). Запись можно осуществлять через:
    * [AWS SDK](../../data-streams/operations/aws-sdk/send.md);
    * [Kafka API](../../data-streams/kafkaapi/auth.md#example);
    * [HTTP API, совместимый с Amazon Kinesis Data Streams](../../data-streams/kinesisapi/methods/putrecord.md).
1. Настройте [триггер для Yandex Data Streams](../../functions/concepts/trigger/data-streams-trigger.md), указав созданный поток данных в качестве источника.

#### Вариант 2: Отправка через Yandex Message Queue {#api-option-ymq}

1. Создайте очередь Yandex Message Queue. Запись сообщений [осуществляется с помощью cURL](../../message-queue/operations/message-queue-send-message.md#curl).
1. Создайте [триггер для Yandex Message Queue](../../functions/concepts/trigger/ymq-trigger.md), указав созданную очередь в качестве источника.

## Форматы сообщений {#message-format}

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

Общее правило: EventRouter доставляет тело события как есть, а при включенном группировании — JSON-массив тел. Триггер всегда оборачивает событие в JSON-объект `{"messages": [...]}`, даже если событие одно.


Источник | Что доставлял EventRouter | Что доставляет триггер
--- | --- | ---
Таймер | Значение поля **Данные** как есть. Если поле пустое, приемник все равно вызывался, но с пустым телом | JSON-объект с полями `event_metadata` и `details`, в которых идентификатор триггера и значение поля **Данные**
Yandex Message Queue | Тело сообщения как есть | JSON-объект с полями `event_metadata` и `details`. Тело сообщения типа `string` находится в `details.message.body`, рядом с ним — идентификатор очереди и атрибуты сообщения
Yandex Data Streams | Запись как есть | JSON-объект с записями из потока данных без дополнительных полей

Точные примеры сообщений приведены в описании каждого типа триггера.

{% note warning %}

Фильтр и шаблон преобразования применяются к каждому сообщению внутри JSON-объекта, а результат снова упаковывается в JSON-объект. Убрать JSON-объект `{"messages": [...]}` шаблоном нельзя, поэтому приемник в любом случае придется научить его разворачивать.

{% endnote %}

Обратите внимание:

* Метаданные события — идентификатор, время создания, атрибуты сообщения — в EventRouter в приемник не попадали. В триггерах они доступны, и их можно использовать, например, для дедупликации по `event_metadata.event_id`.
* Триггер для Yandex Data Streams принимает и отправляет события только в формате JSON.
* Тело события Yandex Message Queue передается строкой независимо от того, что в нем лежит. Если приемник ожидает JSON-объект, тело нужно разобрать с помощью шаблона преобразования или в коде приемника.

### Шаблоны преобразования {#message-format-jq}

Чтобы содержимое элементов JSON-объекта совпадало с тем, что доставлял EventRouter, задайте в приемнике триггера шаблон преобразования.

**Таймер**

Чтобы получить значение поля **Данные**:

```jq
.details.payload
```

Значение придет строкой. Если в поле записан JSON, разберите его:

```jq
.details.payload | fromjson
```

**Yandex Message Queue**

Чтобы получить только тело сообщения:

```jq
.details.message.body
```

Тело придет строкой. Если в очередь пишется JSON, разберите его:

```jq
.details.message.body | fromjson
```

Если в очереди могут оказаться сообщения, не являющиеся корректным JSON, используйте безопасный вариант — он вернет разобранный объект либо исходную строку:

```jq
.details.message.body | fromjson? // .
```

Тело можно дополнить метаданными, которых в EventRouter не было. Например, чтобы передать в приемник тело вместе с идентификатором события:

```jq
{body: (.details.message.body | fromjson), event_id: .event_metadata.event_id}
```

**Yandex Data Streams**

Шаблон не нужен: элементы JSON-объекта — это записи из потока, они совпадают с тем, что доставлял EventRouter. Отличается только JSON-объект.

### Перенос существующих фильтров и шаблонов преобразования {#migrate-existing-jq}

Если в правиле или приемнике EventRouter уже были заданы фильтры и шаблоны преобразования, допишите к ним слева распаковку тела события. Ниже `<выражение>` — это фильтр или шаблон, который был задан в EventRouter.

**Таймер**

```jq
.details.payload | fromjson | <выражение>
```

Например, фильтр `.firstName == "Ivan"` для очереди Yandex Message Queue превращается в:

```jq
.details.message.body | fromjson | .firstName == "Ivan"
```

А шаблон `{name: .firstName, city: .address.city}` — в:

```jq
.details.message.body | fromjson | {name: .firstName, city: .address.city}
```

Если выражение не удалось вычислить (например, тело сообщения не является корректным JSON), событие перенаправляется в Dead Letter Queue приемника, а если она не настроена, теряется. Для триггера с источником Yandex Message Queue Dead Letter Queue недоступна, поэтому такое событие теряется. Для очередей, в которых могут оказаться сообщения произвольного формата, лучше использовать безопасный вариант с `fromjson? // .`.

**Yandex Message Queue**

```jq
.details.message.body | fromjson | <выражение>
```

**Yandex Data Streams**

Выражение переносится без изменений:

```jq
<выражение>
```

## Миграция приемников {#target-migration}

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

Приемник EventRouter | Приемник триггера | Что меняется
--- | --- | ---
Функция | Функция | Настройки группирования задаются на источнике, а не на приемнике
Контейнер | Контейнер | Настройки группирования задаются на источнике, а не на приемнике, нельзя закрепить ревизию контейнера
Рабочий процесс | Рабочий процесс | Настройки группирования задаются на источнике, а не на приемнике
WebSocket-соединения | WebSocket-соединения | Настройки группирования задаются на источнике, а не на приемнике, не поддерживаются повторные вызовы и Dead Letter Queue
Лог-группа | Нет аналога | Нужна [функция-прослойка](#shim)
Поток данных | Нет аналога | Нужна [функция-прослойка](#shim)
Очередь сообщений | Нет аналога | Нужна [функция-прослойка](#shim)

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

### Функции {#target-functions}

Создайте [триггер, вызывающий функцию](../../functions/concepts/trigger/index.md). Идентификатор функции, тег версии и сервисный аккаунт переносятся из приемника без изменений.

Сервисному аккаунту нужна роль `functions.functionInvoker` на функцию, которую вызывает триггер.

Как и в EventRouter, триггер вызывает функцию с параметром строки запроса `?integration=raw`, поэтому способ разбора входных данных в коде функции менять не нужно — меняется только [формат сообщения](#message-format).

### Контейнеры {#target-containers}

Создайте [триггер, вызывающий контейнер](../../serverless-containers/concepts/trigger/index.md). Идентификатор контейнера, путь и сервисный аккаунт переносятся из приемника без изменений.

Сервисному аккаунту нужна роль `serverless-containers.containerInvoker` на контейнер, который вызывает триггер.

{% note warning %}

В приемнике EventRouter можно было указать конкретную ревизию контейнера. В триггере такой настройки нет — всегда вызывается активная ревизия. Если вы закрепляли ревизию, чтобы контролировать момент выкатки новой версии, продумайте замену: например, разделите контейнеры для стабильной и тестовой версий.

{% endnote %}

### Рабочие процессы {#target-workflows}

В приемнике триггера укажите идентификатор рабочего процесса и сервисный аккаунт, от имени которого он будет запускаться. Оба параметра переносятся из приемника без изменений.

Сервисному аккаунту нужна роль `serverless.workflows.executor` на рабочий процесс, который запускает триггер.

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

### WebSocket-соединения {#target-websocket}

Создайте [триггер, отправляющий сообщения в WebSocket-соединения](../../api-gateway/concepts/trigger/index.md). Идентификатор API-шлюза, путь и сервисный аккаунт переносятся из приемника без изменений.

Сервисному аккаунту нужна роль `api-gateway.websocketBroadcaster` на каталог, в котором находится API-шлюз.

{% note warning %}

Для этого типа приемника не поддерживаются повторные вызовы и Dead Letter Queue. Если указать их при создании триггера, ошибки не будет, но настройки не применятся. Если в приемнике EventRouter были настроены повторные попытки или очередь Dead Letter Queue, перенести их не получится.

{% endnote %}

### Лог-группы, потоки данных и очереди сообщений {#shim}

Триггеры не умеют записывать события в лог-группу Cloud Logging, поток Yandex Data Streams или очередь Yandex Message Queue напрямую. Вместо приемника такого типа укажите в триггере приемник с функцией-прослойкой, которая перекладывает события в нужное назначение.

Все функции ниже устроены одинаково: принимают JSON-объект `{"messages": [...]}` и записывают каждый его элемент как отдельную запись. Если в назначение нужно передавать не событие целиком, а только его часть, не меняйте код функции — задайте в приемнике триггера шаблон преобразования. Подробнее в [Шаблоны преобразования](#message-format-jq).

Общее для всех трех функций, приведенных в разделах ниже:

* среда выполнения — `golang123`, точка входа — `index.Handler`;
* вместе с `index.go` загружайте файл `go.mod`. Имя модуля в нем не должно быть `main`. Чтобы зафиксировать версии зависимостей, загрузите еще и `go.sum`, иначе установятся последние;
* в `go.mod` не должно быть строк `go` и `toolchain`. Версия Go в собранном плагине обязана совпадать с версией среды выполнения, и сборщик подставляет ее сам, а эти директивы заставят его взять другую. Функция при этом соберется, но при вызове упадет с ошибкой `fatal error: runtime: no plugin module data`. Строку `toolchain` Go дописывает автоматически при `go get` и `go mod tidy`, поэтому перед загрузкой проверьте файл: в нем должны остаться только `module` и `require`;
* функция вызывается триггером, поэтому в случае ошибки триггер повторит вызов со всем пакетом целиком. Часть событий при этом может быть записана повторно — учитывайте это при обработке. Если источник триггера — очередь сообщений, повторных вызовов не будет;
* сервисному аккаунту, указанному в настройках функции, нужна роль на запись в лог-группу, поток данных или очередь сообщений, подробнее в описании каждой функции.

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

#### Запись в лог-группу {#shim-logs}

Функция пишет каждое событие в стандартный поток вывода. Записи попадают в ту лог-группу, которая указана в [настройках логирования](../../functions/operations/function/logs-write.md) функции. Задайте в них лог-группу, которая была приемником в EventRouter.

Файл `index.go`:

```go
package main

import (
	"bytes"
	"context"
	"encoding/json"
	"fmt"
)

type Request struct {
	Messages []json.RawMessage `json:"messages"`
}

func Handler(ctx context.Context, req *Request) (string, error) {
	var buf bytes.Buffer
	for _, message := range req.Messages {
		buf.Reset()
		if err := json.Compact(&buf, message); err != nil {
			// Событие не является корректным JSON — пишем как есть.
			fmt.Println(string(message))
			continue
		}
		fmt.Println(buf.String())
	}
	return "ok", nil
}
```

Каждая строка, выведенная функцией, становится отдельной записью в лог-группе. Триггер передает JSON-объект с отступами, поэтому событие внутри него занимает несколько строк, и печатать его без предварительного схлопывания нельзя — одно событие превратилось бы в несколько записей. `json.Compact` убирает переносы и отступы.

Что настроить:

* в параметрах функции укажите нужную лог-группу;
* задайте переменную окружения `STRUCTURED_LOGGING` со значением `false`.

{% note warning %}

Без переменной `STRUCTURED_LOGGING=false` однострочная JSON-запись, в которой есть поле `message` или `msg`, будет распознана как [структурированный лог](../../functions/concepts/logs.md#structured-logs). Тогда значение этого поля станет текстом записи, а остальные поля события уедут в `json_payload`. Если события могут содержать поле с таким именем, переменную нужно задать обязательно, иначе записи в лог-группе не будут совпадать с тем, что писал EventRouter.

{% endnote %}

#### Запись в поток данных {#shim-yds}

Функция пишет события в поток данных по протоколу, совместимому с Amazon Kinesis Data Streams, пакетами до 500 записей.

Файл `index.go`:

```go
package main

import (
	"context"
	"encoding/json"
	"fmt"
	"os"
	"time"

	"github.com/aws/aws-sdk-go-v2/aws"
	"github.com/aws/aws-sdk-go-v2/credentials"
	"github.com/aws/aws-sdk-go-v2/service/kinesis"
	"github.com/aws/aws-sdk-go-v2/service/kinesis/types"
)

const (
	endpoint     = "https://yds.serverless.yandexcloud.net"
	region       = "ru-central1"
	maxBatchSize = 500
)

type Request struct {
	Messages []json.RawMessage `json:"messages"`
}

var streamName = os.Getenv("STREAM_NAME")

var client = kinesis.NewFromConfig(aws.Config{
	Region: region,
	Credentials: credentials.NewStaticCredentialsProvider(
		os.Getenv("AWS_ACCESS_KEY_ID"),
		os.Getenv("AWS_SECRET_ACCESS_KEY"),
		"",
	),
}, func(o *kinesis.Options) {
	o.BaseEndpoint = aws.String(endpoint)
})

func Handler(ctx context.Context, req *Request) (string, error) {
	for start := 0; start < len(req.Messages); start += maxBatchSize {
		end := start + maxBatchSize
		if end > len(req.Messages) {
			end = len(req.Messages)
		}

		records := make([]types.PutRecordsRequestEntry, 0, end-start)
		for i, message := range req.Messages[start:end] {
			records = append(records, types.PutRecordsRequestEntry{
				Data:         []byte(message),
				PartitionKey: aws.String(fmt.Sprintf("%d-%d", time.Now().UnixNano(), start+i)),
			})
		}

		out, err := client.PutRecords(ctx, &kinesis.PutRecordsInput{
			StreamName: aws.String(streamName),
			Records:    records,
		})
		if err != nil {
			return "", err
		}
		if failed := aws.ToInt32(out.FailedRecordCount); failed > 0 {
			return "", fmt.Errorf("не удалось записать %d записей", failed)
		}
	}
	return "ok", nil
}
```

Файл `go.mod`:

```
module ydswriter

require (
	github.com/aws/aws-sdk-go-v2 v1.40.1
	github.com/aws/aws-sdk-go-v2/credentials v1.19.10
	github.com/aws/aws-sdk-go-v2/service/kinesis v1.43.1
)
```

Что настроить:

* переменная окружения `STREAM_NAME` — полное имя потока в формате `/kz1/<идентификатор_облака>/<идентификатор_базы_данных>/<имя_потока>`;
* переменные окружения `AWS_ACCESS_KEY_ID` и `AWS_SECRET_ACCESS_KEY` — [статический ключ доступа](../../iam/concepts/authorization/access-key.md) сервисного аккаунта. Секретную часть ключа передавайте через Yandex Lockbox, а не открытым текстом;
* роль `yds.writer` на поток данных для сервисного аккаунта, которому принадлежит ключ.

#### Запись в очередь сообщений {#shim-ymq}

Функция пишет события в очередь по протоколу, совместимому с Amazon SQS, пакетами до 10 сообщений.

Файл `index.go`:

```go
package main

import (
	"context"
	"encoding/json"
	"fmt"
	"os"

	"github.com/aws/aws-sdk-go-v2/aws"
	"github.com/aws/aws-sdk-go-v2/credentials"
	"github.com/aws/aws-sdk-go-v2/service/sqs"
	"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)

const (
	endpoint     = "https://message-queue.api.cloud.yandex.net/"
	region       = "ru-central1"
	maxBatchSize = 10
)

type Request struct {
	Messages []json.RawMessage `json:"messages"`
}

var queueURL = os.Getenv("QUEUE_URL")

var client = sqs.NewFromConfig(aws.Config{
	Region: region,
	Credentials: credentials.NewStaticCredentialsProvider(
		os.Getenv("AWS_ACCESS_KEY_ID"),
		os.Getenv("AWS_SECRET_ACCESS_KEY"),
		"",
	),
}, func(o *sqs.Options) {
	o.BaseEndpoint = aws.String(endpoint)
})

func Handler(ctx context.Context, req *Request) (string, error) {
	for start := 0; start < len(req.Messages); start += maxBatchSize {
		end := start + maxBatchSize
		if end > len(req.Messages) {
			end = len(req.Messages)
		}

		entries := make([]types.SendMessageBatchRequestEntry, 0, end-start)
		for i, message := range req.Messages[start:end] {
			entries = append(entries, types.SendMessageBatchRequestEntry{
				Id:          aws.String(fmt.Sprintf("%d", start+i)),
				MessageBody: aws.String(string(message)),
			})
		}

		out, err := client.SendMessageBatch(ctx, &sqs.SendMessageBatchInput{
			QueueUrl: aws.String(queueURL),
			Entries:  entries,
		})
		if err != nil {
			return "", err
		}
		if len(out.Failed) > 0 {
			return "", fmt.Errorf("не удалось отправить %d сообщений, первая ошибка: %s",
				len(out.Failed), aws.ToString(out.Failed[0].Message))
		}
	}
	return "ok", nil
}
```

Файл `go.mod`:

```
module ymqwriter

require (
	github.com/aws/aws-sdk-go-v2 v1.40.1
	github.com/aws/aws-sdk-go-v2/credentials v1.19.10
	github.com/aws/aws-sdk-go-v2/service/sqs v1.42.21
)
```

Что настроить:

* переменная окружения `QUEUE_URL` — URL очереди. Обратите внимание, что в EventRouter приемник задавался идентификатором очереди в формате ARN, а здесь нужен именно URL;
* переменные окружения `AWS_ACCESS_KEY_ID` и `AWS_SECRET_ACCESS_KEY` — статический ключ доступа сервисного аккаунта;
* роль `ymq.writer` на очередь для сервисного аккаунта, которому принадлежит ключ.