HTTP-расширение для YTsaurus Flow
HTTP-расширение предоставляет асинхронный синк NYT::NFlow::TAsyncHttpSink. Он отправляет каждое входное сообщение в HTTP- или HTTPS-эндпоинт методом POST и ждёт успешного ответа, прежде чем отметить сообщение доставленным.
Синк принимает один входной поток ровно с одной колонкой типа string. Значение не null из колонки, указанной в payload_column, побайтно и без изменений становится телом запроса. Это может быть protobuf, YSON, JSON или произвольная последовательность байтов в значении типа string. Синк не разбирает, не сериализует и не преобразует значение. Значение null пропускается без HTTP-запроса. Заголовки, например Content-Type, задаются в разделе headers параметров синка в Spec.
Код расширения находится здесь.
Настройка синка
Объявите синк в компьютейшене, который создаёт поток с телами запросов. Например:
"computations" = {
"<computation>" = {
"sinks" = {
"http" = {
"sink_class_name" = "NYT::NFlow::TAsyncHttpSink";
"input_stream_ids" = ["requests"];
"parameters" = {
"url" = "https://receiver.example/api/events";
"payload_column" = "payload";
"headers" = {
"Content-Type" = "application/octet-stream";
};
"idempotency_header" = "Idempotency-Key";
"keep_alive" = %true;
"max_redirect_count" = 0;
"max_idle_connections" = 8;
};
};
};
};
};
Значения статических заголовков хранятся в Spec пайплайна в Cypress и доступны пользователям и сервисам с правом чтения Spec, в том числе в дампах спеки. Безопасного механизма для секретных заголовков пока нет. Не помещайте в headers OAuth-токены, API-токены и другие секреты. Не используйте этот синк с эндпоинтами, которым нужны секретные заголовки, если безопасность не обеспечена контролем доступа и внешним механизмом.
Поток requests в этом примере должен иметь следующую схему:
[
{name = "payload"; type = "string"; required = %false;};
]
Все параметры из примера статические. Перед их изменением остановите пайплайн. Параметры повторов динамические. Задайте их в DynamicSpec по пути computations/<computation>/sinks/<sink>/parameters, где <computation> и <sink> — имена из Spec:
"computations" = {
"<computation>" = {
"sinks" = {
"<sink>" = {
"parameters" = {
"request_timeout" = "60s";
"attempt_timeout" = "10s";
"retry_initial_delay" = "1s";
"retry_minimum_delay" = "100ms";
"retry_multiplier" = 2.0;
"retry_maximum_delay" = "30s";
"retry_jitter_ratio" = 0.2;
"max_attempt_count" = 5;
};
};
};
};
};
Обновлённые динамические параметры используются при последующих решениях о повторе, в том числе для уже выполняющегося запроса.
Статическая спека
HTTP-синк поддерживает только доставку at-least-once. Общая схема синков добавляет at_most_once_strategy, но при включении этой стратегии загрузка спеки завершается ошибкой.
Источник: yt/yt/flow/library/cpp/common/registry-inl.h
|
Параметр |
Описание |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
Динамическая спека
Источник: yt/yt/flow/library/cpp/common/registry-inl.h
|
Параметр |
Описание |
|
|
Тип: |
|
|
Тип: TDuration |
|
|
Тип: TDuration |
|
|
Тип: TDuration |
|
|
Тип: TDuration |
|
|
Тип: |
|
|
Тип: TDuration |
|
|
Тип: |
|
|
Тип: |
Доставка и повторы
Синк считает успешным любой числовой статус от 200 до 299. Перед завершением попытки он полностью читает тело ответа, чтобы соединение keep-alive можно было использовать снова. Ответ с другим статусом, сетевая ошибка, тайм-аут попытки или ошибка чтения тела ответа приводят к повтору, пока это позволяют и max_attempt_count, и request_timeout.
Синк проходит не более max_redirect_count перенаправлений во всех попытках одной доставки, сохраняя метод POST, тело и заголовки запроса. Относительное значение Location разрешается относительно URL предыдущего запроса. Синк дочитывает и учитывает каждый ответ редиректа, а также отклоняет редиректы с HTTPS на HTTP. Значение по умолчанию — ноль, поэтому переходы включаются только явно. Задавайте положительный предел, только если всем возможным адресатам можно доверить настроенные заголовки и тело запроса.
request_timeout ограничивает всю доставку вместе с попытками и паузами. Каждая попытка ограничена меньшим из двух интервалов: attempt_timeout и временем, оставшимся до общего дедлайна. Пауза растёт от retry_initial_delay с множителем retry_multiplier, случайно изменяется в пределах retry_jitter_ratio и остаётся между retry_minimum_delay и retry_maximum_delay.
Число повторов конечно. После исчерпания попыток или общего дедлайна доставка окончательно завершается ошибкой на всё время жизни текущих синка и джобы. В этой джобе синк больше не пытается автоматически отправить сообщение методом POST. Сообщение остаётся неподтверждённым, продвижение по входным сообщениям и сорсу этой партиции останавливается, а лаг растёт. При этом обычные эпохи могут продолжать коммитить работу, для которой нет блокирующей доставки.
Повторная отправка возобновляется только после рестарта джобы или воркера либо другого пересоздания синка, когда загружается max_persisted_message_id. Flow не гарантирует автоматический рестарт после исчерпания повторов. Если джоба продолжает работать, после устранения причины перезапустите её штатным способом.
Гарантии и повторная отправка
HTTP-синк обеспечивает доставку at-least-once. Flow записывает завершение только после получения и полного чтения ответа 2xx, поэтому неуспешно доставленное сообщение не теряется. Однако получатель может обработать запрос, после чего соединение оборвётся до того, как Flow запишет успех. При повторе или рестарте воркера то же тело будет отправлено ещё раз.
POST-запросы выполняются параллельно, но для упорядоченного сохранения успешные завершения становятся видны только непрерывным префиксом в порядке отправки. Более поздний ответ 2xx ждёт завершения всех предыдущих доставок. Если первое неподтверждённое сообщение окончательно завершилось ошибкой, подтверждение всех последующих сообщений блокируется на всё время жизни текущих синка и джобы. Поэтому max_persisted_message_id не может перескочить через него или продвинуться с нарушением порядка. Барьер подтверждений не ограничивает параллелизм HTTP-запросов.
Значение null в payload_column подтверждается и пропускается без HTTP-запроса. Оно не блокирует последующие сообщения и не отправляется повторно после сохранения результата.
Получатель должен выполнять операцию идемпотентно или дедуплицировать запросы. По умолчанию синк добавляет заголовок Idempotency-Key, значение которого равно шестнадцатеричному представлению стабильного идентификатора сообщения Flow. При повторах и повторной отправке одного сообщения используется то же значение. Чтобы переименовать заголовок, задайте другое имя в idempotency_header; чтобы отключить его, задайте пустую строку. Статический заголовок с тем же именем из headers не допускается.
Внешний HTTP-вызов не входит в транзакцию эпохи и не покрывается exactly-once гарантией Flow. Внутренний стейт, оффсеты сорсов и сохранение выходных сообщений продолжают подчиняться обычной семантике Flow.
Повторное использование соединений
По умолчанию keep_alive включён. Каждый настроенный HTTP-синк в каждой активной джобе партиции владеет отдельным экземпляром синка, долгоживущими HTTP- и HTTPS-клиентами и отдельным пулом простаивающих соединений для каждого клиента. max_idle_connections задаёт максимальное число соединений, сохраняемых в каждом пуле; значение по умолчанию — восемь, поэтому синк, следующий перенаправлениям через обе схемы, может сохранять до удвоенного числа соединений. Это верно и для TTransformOrderedSourceComputation и TSwiftOrderedSourceComputation: партиции ordered-source соответствует один ключ сорса, поэтому дополнительного множителя по числу ключей внутри джобы нет. Установите keep_alive в %false, если эндпоинт или промежуточный компонент не поддерживает постоянные соединения: тогда ни один клиент этого синка не будет сохранять простаивающие соединения независимо от max_idle_connections. Эти параметры управляют повторным использованием соединений, а не числом попыток.
Каждый экземпляр синка отправляет запросы для всех сообщений своей эпохи одновременно. Число простаивающих соединений ограничивает переиспользование соединений, а не число одновременных запросов: у синка нет параметра для ограничения скорости или параллелизма. Если несколько настроенных синков обращаются к одному эндпоинту, рассчитывайте его параллелизм и число соединений как сумму по всем этим синкам и активным джобам партиций. Вклад каждого синка равен его частоте запросов от одной партиции, умноженной на число активных джоб партиций с этим синком. Ограничивайте параллелизм через входную нагрузку или на стороне получателя.
Метрики
Сенсоры публикуются по пути /sink/async_http_sink/ и содержат тег sink_id с именем синка из Spec:
/sink/async_http_sink/responses{status_code="<code>"}учитывает каждый HTTP-ответ один раз, в том числе успешные статусы и статусы, вызвавшие повтор. В тегеstatus_codeсохраняется числовое значение от эндпоинта./sink/async_http_sink/attempt_failuresучитывает сетевые ошибки и тайм-ауты до получения ответа, а также ошибки чтения тела уже полученного ответа.
Сначала проверьте агрегированный статус пайплайна по пути /sinks/<sink>/async_http_sink: там указана ошибка самой старой доставки, которая сейчас завершается с ошибкой. Статус очищается, только когда не остаётся ни одной доставки с ошибкой. Пока синк находится в состоянии ошибки, соответствующий сенсор /status_profiler/broken{path="/sinks/<sink>/async_http_sink"} равен 1. После этого используйте счётчики, чтобы определить тип сбоя.
Значение null в payload_column пропускается без HTTP-запроса и не увеличивает ни один из счётчиков.
Рост счётчика ответов с кодами вне диапазона 2xx указывает на отказ на уровне приложения. Если растёт attempt_failures, а соответствующий счётчик ответа — нет, проверьте транспорт и тайм-ауты до получения ответа. Ошибка чтения тела ответа увеличивает и attempt_failures, и счётчик полученного статуса, даже если это статус 2xx. Все эти ошибки могут приводить к повторам и последующей повторной отправке сообщения.
Диагностика ошибок
|
Текст ошибки |
Причина и исправление |
|
|
В |
|
|
В |
|
|
В |
|
|
В |
|
|
У синка нет входного потока или в |
|
|
В схеме входного потока нет колонок или их несколько. Оставьте ровно одну колонку с телом запроса. |
|
|
|
|
|
Закончился конечный бюджет повторов. Исправьте эндпоинт, транспорт или тайм-ауты, затем перезапустите джобу или воркер либо иначе пересоздайте синк, чтобы возобновить отправку. В текущей джобе новых POST-запросов автоматически не будет. |
|
|
Пайплайн попал на воркер, собранный до добавления расширения. Обновите все воркеры, на которых может работать пайплайн, до версии со стандартным |