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

Параметр

Описание

at_most_once_strategy

Тип: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyParameters>
Необязательная стратегия доставки at-most-once. Поддержка зависит от коннектора; перед включением проверьте документацию выбранного коннектора.

url

Тип: std::string
Обязательный параметр
Целевой HTTP- или HTTPS-адрес. URL должен содержать имя хоста.

payload_column

Тип: std::string
Обязательный параметр
Имя колонки с телом запроса. Во входном потоке должна быть ровно одна String-колонка с этим именем. Значение не null побайтно и без изменений становится телом запроса: это может быть protobuf, YSON, JSON или произвольная последовательность байтов в значении типа String. Значение null подтверждается и пропускается без HTTP-запроса.

headers

Тип: THashMap<std::string, std::string>
Значение по умолчанию: {}
Статические заголовки HTTP-запроса. Имена должны быть корректными HTTP-токенами. Если получатель определяет формат непрозрачного тела по Content-Type, задайте этот заголовок здесь. Значения хранятся в Spec пайплайна в Cypress и доступны пользователям и сервисам с правом чтения Spec, в том числе в дампах спеки. Безопасного механизма для секретных заголовков пока нет. Не помещайте сюда OAuth-токены, API-токены и другие секреты. Не используйте этот синк с эндпоинтами, которым нужны секретные заголовки, если безопасность не обеспечена контролем доступа и внешним механизмом.

idempotency_header

Тип: std::string
Значение по умолчанию: Idempotency-Key
Имя заголовка с идентификатором сообщения. Значение равно шестнадцатеричному представлению стабильного идентификатора сообщения Flow и не меняется при повторах и повторной отправке. По умолчанию используется Idempotency-Key. Другое корректное имя переименовывает заголовок, а пустая строка отключает его. Статический заголовок с тем же именем не допускается.

keep_alive

Тип: bool
Значение по умолчанию: true
Разрешает повторно использовать HTTP-соединения. Каждый настроенный HTTP-синк в каждой активной джобе партиции владеет отдельными экземпляром синка, клиентом и пулом простаивающих соединений. Если параметр выключен, пул не сохраняет соединения независимо от max_idle_connections.

max_redirect_count

Тип: int
Значение по умолчанию: 0
Максимальное число HTTP-перенаправлений во всех попытках одной доставки с сохранением метода POST, тела и заголовков. Значение по умолчанию — ноль, поэтому переходы включаются только явно. Доверяйте всем возможным адресатам.

max_idle_connections

Тип: int
Значение по умолчанию: 8
Максимальное число простаивающих соединений в пуле каждого HTTP- или HTTPS-клиента настроенного синка и активной джобы партиции, когда включён keep_alive. Значение по умолчанию — восемь на клиент, поэтому синк, следующий перенаправлениям через обе схемы, может сохранять до удвоенного числа соединений. Партиции ordered-source соответствует один ключ сорса, поэтому дополнительного множителя по числу ключей нет. Если несколько синков обращаются к одному эндпоинту, рассчитывайте его ёмкость как сумму по всем настроенным синкам и активным джобам партиций.

Динамическая спека

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

at_most_once_strategy

Тип: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyDynamicParameters>
Значение по умолчанию: {}
Динамические параметры at_most_once_strategy. Поддержка зависит от коннектора; перед настройкой проверьте документацию выбранного коннектора.

request_timeout

Тип: TDuration
Значение по умолчанию: 1m
Общий дедлайн доставки одного сообщения, включая все HTTP-попытки и паузы между ними. Его исчерпание окончательно останавливает синк на время текущей джобы.

attempt_timeout

Тип: TDuration
Значение по умолчанию: 10s
Максимальная длительность одной HTTP-попытки. Попытка дополнительно ограничивается временем, оставшимся до request_timeout; значение не может превышать request_timeout.

retry_initial_delay

Тип: TDuration
Значение по умолчанию: 1s
Пауза перед первым повтором. Значение должно находиться между retry_minimum_delay и retry_maximum_delay.

retry_minimum_delay

Тип: TDuration
Значение по умолчанию: 100ms
Нижняя граница паузы после добавления джиттера. Значение не может превышать retry_initial_delay.

retry_multiplier

Тип: double
Значение по умолчанию: 2.0
Множитель экспоненциального роста паузы. Допустимы значения от 1.01 до 100.0 включительно.

retry_maximum_delay

Тип: TDuration
Значение по умолчанию: 30s
Верхняя граница паузы между попытками. Значение не может быть меньше retry_initial_delay.

retry_jitter_ratio

Тип: double
Значение по умолчанию: 0.2
Случайное относительное отклонение паузы. Допустимы значения от 0.0 до 1.0 включительно.

max_attempt_count

Тип: int
Значение по умолчанию: 5
Максимальное число HTTP-попыток с учётом первой. Его исчерпание окончательно останавливает синк на время текущей джобы. Допустимы значения от 1 до 100 включительно.

Доставка и повторы

Синк считает успешным любой числовой статус от 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. Все эти ошибки могут приводить к повторам и последующей повторной отправке сообщения.

Диагностика ошибок

Текст ошибки

Причина и исправление

URL must use http or https and contain a host

В url указана другая схема или отсутствует хост. Задайте полный URL с хостом и схемой http:// или https://.

header names must be valid HTTP tokens

В headers есть пустой или недопустимый ключ. Используйте корректное имя HTTP-заголовка.

idempotency_header must be empty or a valid HTTP header name

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

idempotency_header must not be hop-by-hop or transport-managed

В idempotency_header указано имя транспортного или сквозного заголовка, например Content-Length или Connection. Используйте имя заголовка уровня приложения.

expects exactly one input stream

У синка нет входного потока или в input_stream_ids указано несколько потоков. Подключите ровно один входной поток.

expects exactly one payload column

В схеме входного потока нет колонок или их несколько. Оставьте ровно одну колонку с телом запроса.

must be the only String column

payload_column не совпадает с именем единственной колонки или эта колонка имеет тип, отличный от string. Приведите имя колонки в соответствие с payload_column, а её тип задайте как string.

Async HTTP POST retry policy exhausted или Async HTTP POST retry deadline exhausted

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

No sink class "NYT::NFlow::TAsyncHttpSink" is registered

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

См. также

Предыдущая
Следующая