Миграция с EventRouter на триггеры
В качестве альтернативы EventRouter вы можете использовать:
- триггеры для вызова функций Cloud Functions;
- триггеры для вызова контейнеров Serverless Containers;
- триггеры для отправки событий в WebSocket-соединения;
- триггеры для запуска рабочих процессов Workflows.
В отличие от EventRouter, где шина — это набор правил и коннекторов, триггер — это один источник и один или несколько приемников.
Что переносится автоматически
Некоторые шины будут перенесены автоматически: каждый коннектор станет источником отдельного триггера, а приемники правил шины — приемниками триггера. Некоторые шины нужно перенести самостоятельно.
Шина не будет перенесена автоматически, если выполняется хотя бы одно из условий:
- в правилах шины заданы фильтр или шаблон преобразования;
- коннектор имеет тип Audit Trails или API EventRouter;
- среди приемников любого правила шины есть поток данных Yandex Data Streams, лог-группа Yandex Cloud Logging или очередь сообщений Yandex Message Queue;
- статус коннектора — не
Запущени неОстановлен; - у коннектора с типом Таймер задан часовой пояс.
Важно
Мы рекомендуем перенести все шины самостоятельно. Автоматический перенос не гарантирует полную идентичность функционирования системы после миграции.
Соответствие сущностей EventRouter и триггеров
| EventRouter | Триггеры | Что учесть при переносе |
|---|---|---|
| Шина | Нет аналога | Триггер связывает источник и приемники напрямую, промежуточной шины нет. |
| Коннектор | Источник триггера | Один коннектор соответствует одному триггеру. |
| Правило | Приемник триггера | Правила принадлежат шине, а не коннектору: каждое событие проходит через все правила шины. Поэтому в каждый триггер попадают приемники всех правил шины. |
| Фильтр правила | Фильтр приемника | В EventRouter фильтр один на все приемники правила, в триггере фильтр задается отдельно для каждого приемника. |
| Приемник | Приемник триггера | Не более 5 приемников для одного триггера. Считаются суммарно по всем правилам шины. |
| Шаблон преобразования приемника | Шаблон преобразования приемника | Применяется к событию другого формата, подробнее в Форматы сообщений. |
| Настройки группирования | Настройки группирования для источника | В 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 и прямая отправка в шину. |
План миграции
Шаг 1. Составьте список ресурсов
Составьте список всех шин, коннекторов и правил, которые нужно перенести:
yc serverless eventrouter bus list
yc serverless eventrouter connector list
yc serverless eventrouter rule list
Команды connector list и rule list выводят все коннекторы и правила каталога — отберите нужные по полю bus_id в выводе. Фильтровать по шине с помощью команды нельзя.
Для каждой шины зафиксируйте ее коннекторы с настройками источников и все ее правила с фильтрами и приемниками. Каждый коннектор станет отдельным триггером, а приемники всех правил шины — приемниками каждого из этих триггеров.
Важно
В EventRouter каждое событие проходит через все правила шины вне зависимости от того, какой коннектор его принес. Поэтому в каждый триггер переносите приемники всех правил шины, а не только тех, которые относятся к этому коннектору. Если в шине было несколько коннекторов, один и тот же набор действий повторится в каждом триггере.
Шаг 2. Проверьте ограничения
Прежде чем создавать триггеры, убедитесь, что ваш сценарий переносится без доработок. Дополнительные работы потребуются, если:
-
источник — API EventRouter или события отправляются в шину напрямую: нужно доработать приложение-отправитель;
-
среди приемников есть очередь сообщений Message Queue, поток данных Data Streams или лог-группа Cloud Logging: нужна функция-прослойка;
-
в правилах шины суммарно больше 5 приемников. Что делать, зависит от источника:
- поток данных Data Streams — заведите в потоке данных дополнительного потребителя и создайте второй триггер с тем же потоком данных. Каждый потребитель получает полную копию событий, так что приемники можно разложить по нескольким триггерам. Указывать в двух триггерах одного потребителя нельзя: тогда они поделят события между собой;
- таймер — создайте второй триггер с тем же расписанием;
- очередь сообщений Message Queue — несколько триггеров на одну очередь создать нельзя, поэтому лишние вызовы придется вынести в функцию-прослойку, которая вызовет остальные приемники;
-
в таймере указан часовой пояс или секунды в cron-выражении;
-
у приемников заданы разные настройки группирования: в триггере они общие для всех приемников, поэтому придется выбрать одни. Автоматический перенос в таких случаях берет минимальные значения по всем приемникам;
-
суммарный размер группы превышает 64 КиБ: в EventRouter лимит 256 КиБ;
-
у приемников настроены повторные вызовы или Dead Letter Queue, а источник — очередь сообщений: для такого триггера ни то, ни другое не поддерживается, повторы настраиваются политикой перенаправления самой очереди, и для этого нужна роль
ymq.admin; -
размер отдельного события превышает 230 КБ.
Важно
Событие больше допустимого размера триггер отбрасывает без повторной попытки и без записи в Dead Letter Queue. Убедитесь, что таких событий в источнике нет.
Убедитесь, что в облаке не более 100 триггеров. Эта квота общая для триггеров Cloud Functions, Serverless Containers и API Gateway, и шина с несколькими коннекторами расходует ее быстро. Чтобы повысить квоту, обратитесь в техническую поддержку
Шаг 3. Подготовьте триггеры к переключению
-
Выдайте сервисным аккаунтам роли, необходимые для работы триггера. Набор ролей зависит от типа источника и приемника, подробнее читайте в описании соответствующего триггера.
-
Создайте пробный триггер, в качестве приемника укажите функцию, которая логирует полученное событие. Источник выбирайте так, чтобы не помешать работающему коннектору:
- поток данных — можно взять продовый поток данных, но обязательно с отдельным потребителем. Тогда пробный триггер получит полную копию событий и не затронет коннектор;
- очередь сообщений — только отдельная тестовая очередь. Триггер для продовой очереди начнет разбирать те же события, и часть событий не дойдет до приемников EventRouter;
- таймер — создайте пробный триггер с тем же расписанием.
-
Убедитесь, что формат событий соответствует ожиданиям, и подготовьте jq-шаблоны и изменения в коде приемников.
Шаг 4. Переключитесь на триггеры
Переключение зависит от типа источника:
| Источник | Порядок действий | Что происходит с событиями |
|---|---|---|
| Yandex Message Queue | Остановите коннектор, дождитесь обработки накопленных событий, создайте триггер | События сохраняются в очереди. Коннектор и триггер могут читать одну очередь одновременно, но тогда они поделят события между собой, поэтому проверить два контура параллельно на одной очереди нельзя — для проверки нужна вторая очередь. |
| Yandex Data Streams | Остановите коннектор, создайте триггер с тем же потребителем | События сохраняются в потоке. Если создать триггер с новым потребителем, чтение начнется не с того же места. |
| Таймер | Остановите коннектор, создайте триггер | Одно срабатывание может быть пропущено или продублировано. |
| API EventRouter, прямая отправка в шину | Отправьте из приложения событие в очередь сообщений или поток данных, создайте триггер | События, отправленные в шину после остановки, теряются. |
Триггеры, как и EventRouter, гарантируют доставку At least once, поэтому во время переключения возможны повторные вызовы. Убедитесь, что обработчики идемпотентны.
Шаг 5. Проверьте работу
- Отправьте тестовое событие и убедитесь, что приемник вызван.
- Сравните метрики вызовов приемников до и после переключения.
- Если в настройках приемника указана Dead Letter Queue, проверьте, что она не пополняется. Для триггера с источником Yandex Message Queue такой проверки нет — смотрите DLQ, указанную в настройках политики перенаправления очереди.
Шаг 6. Удалите ресурсы EventRouter
Снимите защиту от удаления, если она включена, и удалите правила, коннекторы и шины. После этого отзовите роли, которые выдавались сервисным аккаунтам только для работы EventRouter.
Пример миграции шины
Ниже разобран типовой случай: шина с одним коннектором и двумя правилами превращается в один триггер с двумя приемниками.
Что было в EventRouter
К шине orders-bus подключен один коннектор orders-queue с источником Yandex Message Queue: очередь orders. События в очереди выглядят так:
{"orderId": "1234", "status": "new", "amount": 500}
К шине привязаны два правила, у обоих приемников задано группирование по 10 событий или 5 секунд:
| Правило | Фильтр | Приемник |
|---|---|---|
process-orders |
.status == "new" |
Контейнер order-processor |
notify-orders |
Не задан | Функция order-notifier |
Что получится в триггерах
Один триггер с источником Yandex Message Queue и двумя приемниками:
| Было | Стало |
|---|---|
Коннектор orders-queue |
Источник триггера: очередь orders |
| Настройки группирования у каждого приемника | Одни настройки группирования на источнике: 10 сообщений или 5 секунд |
Правило process-orders |
Приемник 1: вызов контейнера order-processor с фильтром |
Правило notify-orders |
Приемник 2: вызов функции order-notifier |
Настройки группирования в EventRouter указываются для приемника, поэтому у разных приемников они могли отличаться. В триггере они общие для всех приемников, и при переносе нужно выбрать одно значение. При автоматическом переносе будут использоваться минимальные значения по всем приемникам.
Если бы приемником одного из правил была лог-группа, очередь сообщений или поток данных, вместо приемника понадобилась бы функция-прослойка — прямых аналогов у этих приемников нет.
Шаг 1. Опишите действия
Фильтр правила process-orders нельзя перенести дословно: он был написан для тела события, а приемник триггера получает JSON-объект с событием. По рецепту переноса допишите слева распаковку тела. Тем же выражением задайте шаблон, чтобы контейнер получал внутри JSON-объекта тела событий, а не служебную обертку сообщения очереди.
Приемник 1 — вызов контейнера.
{
"invokeContainer": {
"containerId": "<идентификатор_контейнера_order-processor>",
"serviceAccountId": "<идентификатор_сервисного_аккаунта>"
},
"filter": {"jq": ".details.message.body | fromjson | .status == \"new\""},
"transformer": {"jq": ".details.message.body | fromjson"}
}
Приемник 2 — вызов функции. У правила notify-orders фильтра не было, поэтому и в приемнике его нет.
{
"invokeFunction": {
"functionId": "<идентификатор_функции_order-notifier>",
"serviceAccountId": "<идентификатор_сервисного_аккаунта>"
},
"transformer": {"jq": ".details.message.body | fromjson"}
}
Ни в одном из приемников нет ни повторных вызовов, ни Dead Letter Queue: источник — очередь сообщений, а для такого триггера они не поддерживаются. Если указать для приемника retryPolicy или deadLetter, создание триггера завершится ошибкой. Повторную обработку настраивайте с помощью политики перенаправления самой очереди — для этого нужна роль ymq.admin.
Шаг 2. Создайте триггер
Сохраните описания приемников в файлы action-1.json и action-2.json и создайте триггер:
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.
Важно
Обязательно указывайте v2 в пути команды. Без него вызывается устаревшая группа команд yc serverless trigger v1, которая пока остается вариантом по умолчанию и не поддерживает --action, фильтры, шаблоны и выбор потребителя. Попытка выполнить команду без v2 завершится ошибкой unknown flag: --action.
Триггер с несколькими приемниками, фильтрами и шаблонами нельзя создать отдельными параметрами вида --invoke-function-id — они задают один приемник без дополнительных настроек. Используйте консоль управления, yc serverless trigger v2, API v2 или Terraform.
На каждый приемник приходится один параметр --action, в триггере может быть не больше пяти приемников.
Что изменится для приемников
Группирование было включено и раньше, поэтому контейнер order-processor уже получал не отдельное событие, а JSON-массив тел:
[
{"orderId": "1234", "status": "new", "amount": 500}
]
Он продолжит получать только события со статусом new, но теперь массив будет лежать в JSON-объекте по ключу messages:
{
"messages": [
{"orderId": "1234", "status": "new", "amount": 500}
]
}
Код контейнера нужно научить разворачивать JSON-объект: убрать его шаблоном нельзя.
То же касается функции order-notifier: она получит все события пакетом, в JSON-объекте, и по-прежнему без фильтрации.
Если источник — поток данных
Для коннектора с источником Yandex Data Streams порядок тот же, с двумя отличиями:
- при создании триггера укажите того же потребителя, который был задан в коннекторе, иначе чтение начнется не с того же места;
- фильтр и шаблон переносятся без изменений — элементы JSON-объекта совпадают с записями потока, распаковывать тело не нужно. Фильтр правила остается выражением
.status == "new", а шаблон не нужен.
Миграция источников (коннекторов)
Таймер
Создайте таймер.
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 — субботе.
Важно
Триггеры не поддерживают указание секунд в cron-выражении. Минимальная единица измерения — 1 минута. Если сценарий требует более частого срабатывания, пересмотрите логику работы приложения.
Важно
В триггерах нельзя задать часовой пояс, время в cron-выражении всегда указывается по UTC+0. Если в коннекторе был задан другой часовой пояс, пересчитайте время в расписании самостоятельно. Учтите, что при таком пересчете расписание перестанет автоматически учитывать переход на летнее и зимнее время, если он есть в вашем часовом поясе.
Yandex Message Queue
Создайте триггер для Yandex Message Queue.
Формат сообщения от триггера для Message Queue отличается от формата в EventRouter. Подробнее в Форматы сообщений. Рекомендуем использовать шаблон преобразования в настройках триггера или изменить конфигурацию вызываемых ресурсов для адаптации под ваши задачи. Например, для получения только тела сообщения укажите в настройках триггера шаблон .details.message.body.
Важно
У приемников такого триггера не поддерживаются повторные попытки отправки и Dead Letter Queue. Повторную обработку настраивайте с помощью политики перенаправления самой очереди — для этого нужна роль ymq.admin.
Из настроек коннектора переносится только таймаут видимости сообщения. У размера группы при чтении из очереди и таймаута опроса аналогов в триггерах нет.
Обратите внимание, что в качестве источника триггера задается ARN очереди — так же, как в приемнике EventRouter. URL очереди понадобится только функции-прослойке, если очередь была еще и приемником.
Yandex Data Streams
Создайте триггер для Yandex Data Streams.
Чтобы избежать повторной обработки событий, убедитесь, что при создании триггера указан тот же потребитель, который был настроен в коннекторе. В устаревшей группе команд yc serverless trigger такого поля нет: сервис заведет собственного потребителя с именем по идентификатору триггера, и позиция чтения потеряется. Указывайте потребителя через yc serverless trigger v2 create yds или API v2.
Содержимое записей триггер передает без изменений, но оборачивает их в JSON-объект {"messages": [...]}. Подробнее в Форматы сообщений.
Audit Trails
Прямого триггера для событий Audit Trails в текущей реализации не предусмотрено. Если вы используете EventRouter для обработки аудитных событий сервисов Container Registry или Object Storage, подойдут триггер для Container Registry или триггер для Object Storage.
Для других сценариев необходимо:
- Экспортировать события в поток данных.
- Настроить интеграцию с Audit Trails, создав трейл и указав созданный поток данных как объект назначения.
- Создать триггер для Yandex Data Streams, указав созданный поток данных как источник.
API EventRouter и прямая отправка в шину
Прямого аналога в триггерах не предусмотрено ни для одного из способов отправки пользовательских событий в 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
- Создайте поток данных Yandex Data Streams. Запись можно осуществлять через:
- Настройте триггер для Yandex Data Streams, указав созданный поток данных в качестве источника.
Вариант 2: Отправка через Yandex Message Queue
- Создайте очередь Yandex Message Queue. Запись сообщений осуществляется с помощью cURL.
- Создайте триггер для Yandex Message Queue, указав созданную очередь в качестве источника.
Форматы сообщений
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-объект с записями из потока данных без дополнительных полей |
Точные примеры сообщений приведены в описании каждого типа триггера.
Важно
Фильтр и шаблон преобразования применяются к каждому сообщению внутри JSON-объекта, а результат снова упаковывается в JSON-объект. Убрать JSON-объект {"messages": [...]} шаблоном нельзя, поэтому приемник в любом случае придется научить его разворачивать.
Обратите внимание:
- Метаданные события — идентификатор, время создания, атрибуты сообщения — в EventRouter в приемник не попадали. В триггерах они доступны, и их можно использовать, например, для дедупликации по
event_metadata.event_id. - Триггер для Yandex Data Streams принимает и отправляет события только в формате JSON.
- Тело события Yandex Message Queue передается строкой независимо от того, что в нем лежит. Если приемник ожидает JSON-объект, тело нужно разобрать с помощью шаблона преобразования или в коде приемника.
Шаблоны преобразования
Чтобы содержимое элементов JSON-объекта совпадало с тем, что доставлял EventRouter, задайте в приемнике триггера шаблон преобразования.
Таймер
Чтобы получить значение поля Данные:
.details.payload
Значение придет строкой. Если в поле записан JSON, разберите его:
.details.payload | fromjson
Yandex Message Queue
Чтобы получить только тело сообщения:
.details.message.body
Тело придет строкой. Если в очередь пишется JSON, разберите его:
.details.message.body | fromjson
Если в очереди могут оказаться сообщения, не являющиеся корректным JSON, используйте безопасный вариант — он вернет разобранный объект либо исходную строку:
.details.message.body | fromjson? // .
Тело можно дополнить метаданными, которых в EventRouter не было. Например, чтобы передать в приемник тело вместе с идентификатором события:
{body: (.details.message.body | fromjson), event_id: .event_metadata.event_id}
Yandex Data Streams
Шаблон не нужен: элементы JSON-объекта — это записи из потока, они совпадают с тем, что доставлял EventRouter. Отличается только JSON-объект.
Перенос существующих фильтров и шаблонов преобразования
Если в правиле или приемнике EventRouter уже были заданы фильтры и шаблоны преобразования, допишите к ним слева распаковку тела события. Ниже <выражение> — это фильтр или шаблон, который был задан в EventRouter.
Таймер
.details.payload | fromjson | <выражение>
Например, фильтр .firstName == "Ivan" для очереди Yandex Message Queue превращается в:
.details.message.body | fromjson | .firstName == "Ivan"
А шаблон {name: .firstName, city: .address.city} — в:
.details.message.body | fromjson | {name: .firstName, city: .address.city}
Если выражение не удалось вычислить (например, тело сообщения не является корректным JSON), событие перенаправляется в Dead Letter Queue приемника, а если она не настроена, теряется. Для триггера с источником Yandex Message Queue Dead Letter Queue недоступна, поэтому такое событие теряется. Для очередей, в которых могут оказаться сообщения произвольного формата, лучше использовать безопасный вариант с fromjson? // ..
Yandex Message Queue
.details.message.body | fromjson | <выражение>
Yandex Data Streams
Выражение переносится без изменений:
<выражение>
Миграция приемников
Один триггер поддерживает до пяти приемников. Для каждого приемника можно указать шаблон преобразования или фильтр, поэтому приемники всех правил шины становятся приемниками одного триггера. Типы приемников можно комбинировать: один триггер может одновременно вызывать функции, контейнеры и рабочие процессы и отправлять сообщения в WebSocket-соединения.
| Приемник EventRouter | Приемник триггера | Что меняется |
|---|---|---|
| Функция | Функция | Настройки группирования задаются на источнике, а не на приемнике |
| Контейнер | Контейнер | Настройки группирования задаются на источнике, а не на приемнике, нельзя закрепить ревизию контейнера |
| Рабочий процесс | Рабочий процесс | Настройки группирования задаются на источнике, а не на приемнике |
| WebSocket-соединения | WebSocket-соединения | Настройки группирования задаются на источнике, а не на приемнике, не поддерживаются повторные вызовы и Dead Letter Queue |
| Лог-группа | Нет аналога | Нужна функция-прослойка |
| Поток данных | Нет аналога | Нужна функция-прослойка |
| Очередь сообщений | Нет аналога | Нужна функция-прослойка |
Во всех случаях сервисный аккаунт, от имени которого вызывается приемник, переносится без изменений.
Функции
Создайте триггер, вызывающий функцию. Идентификатор функции, тег версии и сервисный аккаунт переносятся из приемника без изменений.
Сервисному аккаунту нужна роль functions.functionInvoker на функцию, которую вызывает триггер.
Как и в EventRouter, триггер вызывает функцию с параметром строки запроса ?integration=raw, поэтому способ разбора входных данных в коде функции менять не нужно — меняется только формат сообщения.
Контейнеры
Создайте триггер, вызывающий контейнер. Идентификатор контейнера, путь и сервисный аккаунт переносятся из приемника без изменений.
Сервисному аккаунту нужна роль serverless-containers.containerInvoker на контейнер, который вызывает триггер.
Важно
В приемнике EventRouter можно было указать конкретную ревизию контейнера. В триггере такой настройки нет — всегда вызывается активная ревизия. Если вы закрепляли ревизию, чтобы контролировать момент выкатки новой версии, продумайте замену: например, разделите контейнеры для стабильной и тестовой версий.
Рабочие процессы
В приемнике триггера укажите идентификатор рабочего процесса и сервисный аккаунт, от имени которого он будет запускаться. Оба параметра переносятся из приемника без изменений.
Сервисному аккаунту нужна роль serverless.workflows.executor на рабочий процесс, который запускает триггер.
Входными данными запуска становится сообщение в том виде, в котором его доставил триггер. Если на источнике настроено группирование, один запуск получает сразу пакет сообщений.
WebSocket-соединения
Создайте триггер, отправляющий сообщения в WebSocket-соединения. Идентификатор API-шлюза, путь и сервисный аккаунт переносятся из приемника без изменений.
Сервисному аккаунту нужна роль api-gateway.websocketBroadcaster на каталог, в котором находится API-шлюз.
Важно
Для этого типа приемника не поддерживаются повторные вызовы и Dead Letter Queue. Если указать их при создании триггера, ошибки не будет, но настройки не применятся. Если в приемнике EventRouter были настроены повторные попытки или очередь Dead Letter Queue, перенести их не получится.
Лог-группы, потоки данных и очереди сообщений
Триггеры не умеют записывать события в лог-группу Cloud Logging, поток Yandex Data Streams или очередь Yandex Message Queue напрямую. Вместо приемника такого типа укажите в триггере приемник с функцией-прослойкой, которая перекладывает события в нужное назначение.
Все функции ниже устроены одинаково: принимают JSON-объект {"messages": [...]} и записывают каждый его элемент как отдельную запись. Если в назначение нужно передавать не событие целиком, а только его часть, не меняйте код функции — задайте в приемнике триггера шаблон преобразования. Подробнее в Шаблоны преобразования.
Общее для всех трех функций, приведенных в разделах ниже:
- среда выполнения —
golang123, точка входа —index.Handler; - вместе с
index.goзагружайте файлgo.mod. Имя модуля в нем не должно бытьmain. Чтобы зафиксировать версии зависимостей, загрузите еще иgo.sum, иначе установятся последние; - в
go.modне должно быть строкgoиtoolchain. Версия Go в собранном плагине обязана совпадать с версией среды выполнения, и сборщик подставляет ее сам, а эти директивы заставят его взять другую. Функция при этом соберется, но при вызове упадет с ошибкойfatal error: runtime: no plugin module data. СтрокуtoolchainGo дописывает автоматически приgo getиgo mod tidy, поэтому перед загрузкой проверьте файл: в нем должны остаться толькоmoduleиrequire; - функция вызывается триггером, поэтому в случае ошибки триггер повторит вызов со всем пакетом целиком. Часть событий при этом может быть записана повторно — учитывайте это при обработке. Если источник триггера — очередь сообщений, повторных вызовов не будет;
- сервисному аккаунту, указанному в настройках функции, нужна роль на запись в лог-группу, поток данных или очередь сообщений, подробнее в описании каждой функции.
Учитывайте, что прослойка меняет модель эксплуатации:
- вместо декларативной доставки появляется код, который нужно сопровождать;
- для записи в очередь и поток данных требуется статический ключ доступа сервисного аккаунта вместо управляемых прав доступа;
- вызовы функции тарифицируются.
Запись в лог-группу
Функция пишет каждое событие в стандартный поток вывода. Записи попадают в ту лог-группу, которая указана в настройках логирования функции. Задайте в них лог-группу, которая была приемником в EventRouter.
Файл index.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.
Важно
Без переменной STRUCTURED_LOGGING=false однострочная JSON-запись, в которой есть поле message или msg, будет распознана как структурированный лог. Тогда значение этого поля станет текстом записи, а остальные поля события уедут в json_payload. Если события могут содержать поле с таким именем, переменную нужно задать обязательно, иначе записи в лог-группе не будут совпадать с тем, что писал EventRouter.
Запись в поток данных
Функция пишет события в поток данных по протоколу, совместимому с Amazon Kinesis Data Streams, пакетами до 500 записей.
Файл index.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— статический ключ доступа сервисного аккаунта. Секретную часть ключа передавайте через Yandex Lockbox, а не открытым текстом; - роль
yds.writerна поток данных для сервисного аккаунта, которому принадлежит ключ.
Запись в очередь сообщений
Функция пишет события в очередь по протоколу, совместимому с Amazon SQS, пакетами до 10 сообщений.
Файл index.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на очередь для сервисного аккаунта, которому принадлежит ключ.