Гарантии обработки в YTsaurus Flow
Обзор
Flow по умолчанию обеспечивает exactly-once семантику обработки событий. Это означает:
- Потребители получают результат обработки каждого входного сообщения ровно один раз, без потерь и дубликатов.
- Стейт обновляется атомарно вместе с обработкой сообщения.
При необходимости семантику можно ослабить до at-least-once или at-most-once — например, ради снижения задержек или потребления ресурсов пайплайном.
Exactly-once семантика
Как это работает
Exactly-once обеспечивается тремя механизмами, работающими совместно.
1. Lease-транзакция
Lease — мастер-транзакция, которой владеет контроллер; он создаёт её для каждой джобы. Пользователь не создаёт Lease и не управляет её временем жизни. Контроллер поддерживает транзакцию активной периодическими пингами и проверяет состояние всех Lease в каждом цикле планирования.
Каждая обычная транзакция, в которой джоба коммитит данные эпохи, указывает Lease этой джобы как пререквизит в prerequisite_transaction_ids. Для успешного коммита Lease должна оставаться активной. Lease завершается одним из двух способов. Контроллер отменяет её, когда удаляет джобу: при ребалансировке партиций, после сбоя джобы или потери воркера, а также при остановке, паузе или завершении пайплайна. Lease также может истечь сама, если контроллер прекращает поддерживать её пингами. В обоих случаях срабатывает один и тот же барьер.
После отмены или истечения Lease коммиты с этим пререквизитом отклоняются, поэтому зомби-джоба не может записать данные после того, как контроллер перестал считать её владельцем партиции. Даже при отсутствии данных джоба периодически коммитит почти пустую транзакцию с тем же пререквизитом, чтобы обнаружить потерю Lease.
Если Lease истекла, контроллер замечает это в очередном цикле планирования и удаляет соответствующую джобу. Пока пайплайн остаётся активным и планирование продолжается, после отмены или истечения Lease контроллер удаляет старую джобу, если она ещё существует, планирует необходимые новые джобы для затронутых текущих партиций и создаёт для каждой новой джобы новую Lease. Новые джобы восстанавливают состояние из последней закоммиченной эпохи, возобновляют обработку затронутых текущих партиций и могут повторно обработать только незакоммиченную работу. При завершении, остановке или паузе пайплайна контроллер вместо этого удаляет джобы, отменяет их Lease и не планирует замену.
Lease служит барьером владения перед коммитом. Сама по себе она не выполняет дедупликацию входных сообщений, не обеспечивает атомарность выходных данных или атомарность состояния и не распространяется на произвольные внешние побочные эффекты.
2. Дедупликация входных сообщений
- Внутренние потоки (
input): дедупликация поmessage_idс помощью внутренней таблицыinput_messagesили её компактного вариантаcompact_input_messages. Компактный вариант выбирается автоматически, когда не включёнexperimental_enable_non_uint_key(партиционирование идёт исключительно поuint64-хеш-колонке); поведение можно переопределить параметромuse_compact_input_messages. Старыеmessage_idудаляются из этой таблицы с использованием SystemWatermark. - Источники (
source): дедупликация по оффсетам — оффсеты консьюмера продвигаются только после успешного коммита эпохи. - Таймеры: результат обработки таймера коммитится в той же транзакции, что и удаление этого таймера.
3. Атомарность выходных данных
Способ доставки выходных сообщений зависит от типа синка:
- Синхронные синки (например,
TSyncQueueSink) — записывают данные в целевую систему в той же транзакции эпохи. Запись и подтверждение атомарны. - Асинхронные синки (например,
TAsyncQueueSink) — синхронно сохраняют сообщения в YTsaurus в таблицеoutput_messages, а отправляют сообщения в целевую систему уже после коммита эпохи, используя пару(producer_id, seqNo)для дедупликации на стороне приёмника (Queue API). Если отправка не удалась, сообщения изoutput_messagesбудут отправлены повторно с тем жеseqNo. По сути реализуется паттерн transactional outbox с дедупликацией в целевой системе, если она это поддерживает.
Что требуется от разработчика
Exactly-once гарантии Flow распространяются на внутреннее состояние пайплайна и встроенные механизмы доставки. Однако пользовательский код может нарушить эти гарантии, если не соблюдать следующие правила.
Детерминизм вычислений
Требования к детерминизму зависят от типа компьютейшена.
Для Swift-компьютейшенов (TSwiftMapComputation, TSwiftOrderedSourceComputation) детерминизм строго обязателен. В Swift-компьютейшенах результаты вычислений не материализуются между эпохами и могут быть пересчитаны при рестарте. Если разные попытки вычисления одного и того же сообщения дают разные результаты, возможно смешение результатов от разных попыток, потеря или дублирование данных.
Для обычных компьютейшенов (TTransformComputation, TTransformOrderedSourceComputation) недетерминированный код технически допустим: из всех попыток обработки (возникших из-за рестартов джоб) атомарно применяется ровно одна. Смешения результатов не происходит. Однако если код производит разные выходные данные при повторном запуске, это может привести к неожиданному поведению — особенно при наличии побочных эффектов (вроде записи во внешние системы) или при обновлении стейта на основе недетерминированных значений.
Что нарушает детерминизм
- Использование текущего времени (
Now(),time.time()) в бизнес-логике. - Генерация случайных чисел.
- Обращение к внешним сервисам, результат которых может меняться между вызовами.
Побочные эффекты
Любое действие за пределами Flow — HTTP-запрос, запись в файл, отправка уведомления — не покрывается exactly-once гарантиями. При рестарте джобы такое действие может быть выполнено повторно.
Если побочные эффекты неизбежны, разработчик должен самостоятельно обеспечить их идемпотентность (например, используя уникальный идентификатор сообщения как ключ идемпотентности).
Внешний стейт
Обновления внутреннего стейта Flow (YSON-state, External State через IExternalStateManager) атомарны и безопасны. Однако если компьютейшен модифицирует данные в сторонней системе (базе данных, кеше), эти изменения не участвуют в транзакции эпохи и могут быть продублированы.
Что Flow гарантирует автоматически
- Внутренние потоки между компьютейшенами — exactly-once.
- Обновления стейта — атомарны с обработкой сообщения.
- Обработка и удаление таймера — в одной транзакции.
- Продвижение оффсетов сорсов — только после успешного коммита эпохи.
At-least-once семантика
At-least-once гарантирует, что каждое сообщение будет обработано хотя бы один раз, но допускает повторную обработку. Это может быть полезно, когда бизнес-логика естественно идемпотентна, а стоимость дедупликации неоправданно высока.
Режим processing_mode: at_least_once_consistent
Параметр processing_mode задаётся на уровне TTransformComputation в спеке. Значение по умолчанию — exactly_once.
"computations" = {
"MyComputation" = {
"processing_mode" = "at_least_once_consistent";
...
};
};
Что меняется: отключается дедупликация входных сообщений по message_id. Сообщения могут быть обработаны повторно после рестарта джобы.
Что не меняется: стейт остаётся консистентным (отсюда слово «consistent» в названии режима). Транзакционность эпохи, Lease-механизм, атомарность обновления стейта — всё это продолжает работать.
Когда использовать: если накладные расходы на дедупликацию существенны, а бизнес-логика идемпотентна (например, insert_or_assign вместо increment).
Батчинг в Swift-компьютейшенах: allow_batching_with_relaxed_guarantees
Параметр allow_batching_with_relaxed_guarantees задаётся в блоке parameters компьютейшена типа TSwiftMapComputation. Значение по умолчанию — %false.
"computations" = {
"MyBatcher" = {
"parameters" = {
"allow_batching_with_relaxed_guarantees" = %true;
};
...
};
};
Что меняется: одно выходное сообщение может склеивать несколько входных (батчинг) — например, чтобы свернуть множество мелких сообщений в одно крупное и снизить нагрузку по числу сообщений на партицию ниже по конвейеру. Склеенное сообщение получает MessageId, детерминированно выведенный из набора MessageId родителей, а входное сообщение считается обработанным только после доставки всех его детей.
Что это означает для гарантий: при рестарте границы склейки могут отличаться. Реплей, воспроизводящий склейку с тем же составом родителей, даёт тот же MessageId и дедуплицируется; но если родитель попал в склейку другого состава, её MessageId будет другим — нижестоящие компьютейшены должны быть готовы видеть содержимое каждого родителя более одного раза (at-least-once). Кроме того, порядок по MessageId в рамках одного ключа на стороне нижестоящих компьютейшенов не сохраняется: логика, полагающаяся на упорядоченность по MessageId внутри ключа, должна быть переписана с учётом этого.
Когда использовать: когда нужно укрупнить поток сообщений (батчинг/свёртка), а нижестоящая логика идемпотентна и не зависит от порядка MessageId внутри ключа. С выключенным флагом (по умолчанию) каждое выходное сообщение имеет ровно одного родителя и Swift-компьютейшен сохраняет exactly-once.
At-most-once семантика
At-most-once гарантирует, что каждое сообщение будет обработано не более одного раза, но допускает потерю сообщений.
Эта семантика достигается настройкой at_most_once_strategy на асинхронных синках, которые её поддерживают.
"sinks" = {
"my_sink" = {
"at_most_once_strategy" = {
"enabled" = %true;
"total_queue_bytes_limit" = 104857600; // 100 МБ
"suspend_destruction_duration" = 60000; // мс
};
...
};
};
Как работает: сообщения помещаются во внутреннюю очередь. Если очередь заполнена (total_queue_bytes_limit), новые сообщения отбрасываются без ошибки. Дедупликация и гарантии доставки отсутствуют.
Когда использовать: телеметрия, метрики, best-effort уведомления — сценарии, где потеря части сообщений допустима, а пропускная способность важнее надёжности.
Гарантии порядка
Flow не гарантирует глобальный порядок обработки сообщений. Однако существуют гарантии для производных сообщений с совпадающими ключами в рамках одной цепочки lineage.
Подробности — в разделе Порядок обработки сообщений.
Влияние коннекторов на гарантии
Каждый коннектор имеет свой набор гарантий, зависящий от типа подключения (source/sink) и варианта синка (sync/async/at-least-once).
Queue (QYT)
- Source: exactly-once — оффсеты консьюмера продвигаются только после коммита эпохи.
- Синхронный синк (
TSyncQueueSink): exactly-once — запись в очередь происходит в основной транзакции эпохи. Работает только на основном кластере процессинга. - Асинхронный синк (
TAsyncQueueSink): exactly-once — сообщения сохраняются вoutput_messages, затем доставляются сproducer_id+seqNo. Queue API дедуплицирует повторы. Работает кросс-кластерно.
Подробнее — Queue.
Static Table
- Source: exactly-once — дедупликация по диапазонам чтения.
- Синк (
TArrivalOrderTableSink): exactly-once — таблица и прогресс коммитятся одной master-транзакцией; frontier партиции дедуплицирует replay, поэтому при частично покрытом replay записывается только непокрытый хвост без перезапуска job; callback доставки вызывается только после этого внешнего коммита и следующего коммита Flow.
Подробнее — Static Table.
HTTP
- Source: отсутствует.
- Асинхронный синк (
TAsyncHttpSink): at-least-once — POST-запросы выполняются вне транзакции эпохи, поэтому получатель должен выполнять операцию идемпотентно или дедуплицировать стабильный идентификатор сообщения из заголовкаIdempotency-Key, который используется по умолчанию. При включенииat_most_once_strategyзагрузка спеки завершается ошибкой.
Подробнее — HTTP-расширение.
ClickHouse
ClickHouse-расширение предоставляет три уровня гарантий в четырёх классах синков; выбор напрямую влияет на гарантии обработки сообщений:
- Exactly-once синки (
TClickHouseBatchingSink, рекомендованный по умолчанию, иTShardedClickHouseBatchingSink): сообщения сохраняются вoutput_messagesи детерминированно группируются в батчи по одним и тем же границамMessageId. В нешардированном синкеinsert_deduplication_tokenравен максимальномуMessageIdбатча. При использованииshard_hostsкаждый шард получает этот токен батча с суффиксом:<имя шарда>. Байт-идентичный реплей несёт тот же токен дедупликации даже послеgroup_by-репартиционирования, поэтому exactly-once переживает репартиционирование. ClickHouse дедуплицирует повтор по токену. - At-least-once синк (
TAtLeastOnceClickHouseSink): батч эпохи пишется синхронно в ходе обработки без промежуточного хранения и без токена дедупликации. Ниже задержка и нагрузка на YTsaurus. При ошибке записи эпоха не коммитится, поэтому сообщения не теряются. Известные временные ошибки повторяются без ограничения числа попыток, а для неклассифицированных ошибок выполняется не болееmax_insert_attemptsпопыток вставки всего (по умолчанию 10). Постоянные ошибки сразу приводят к отказу. При сбое после успешной записи, но до коммита эпохи, тот же батч запишется повторно, поэтому возможны дубликаты. - At-most-once синк (
TAtMostOnceClickHouseSink): гарантия действует в обоих режимах стратегии. При статическом значении по умолчаниюat_most_once_strategy.enabled = falseожидающие сообщения остаются вoutput_messagesи доставляются по порядку, поэтому ошибку подключения или запуска сессии до началаINSERTможно повторить. Значениеtrueвключает независимую отправку каждого сообщения через ограниченную очередь в памяти без ожидания доставки. Её динамический лимит задаётat_most_once_strategy.total_queue_bytes_limit; при отбрасывании сообщений из-за переполнения синк пишет предупреждение. В обоих режимах после началаINSERTсинк подтверждает сообщение с ошибкой вместо повторной вставки. Каждое сообщение даёт отдельныйINSERT, поэтому, чтобы не создавать множество мелких частей данных в ClickHouse, поток стоит укрупнять выше по пайплайну.
Exactly-once требует один из движков таблиц ClickHouse, дедуплицирующих вставляемые блоки (см. Data Replication и SharedMergeTree в документации ClickHouse):
ReplicatedMergeTreeReplicatedReplacingMergeTreeReplicatedSummingMergeTreeReplicatedAggregatingMergeTreeSharedMergeTree
Каждый экземпляр синка создаёт сессию записи отложенно при первой записи и сохраняет её до своей остановки. Полная проверка метаданных целевой таблицы выполняется при запуске этой сессии. При динамической перенастройке клиенты пересоздаются после изменения таймаута записи, а при изменении связанных параметров обновляется проверка окна дедупликации. Движок, схема и идентификатор репликации повторно не проверяются. После изменения целевой таблицы приостановите и снова запустите пайплайн, чтобы новые экземпляры синка проверили её. Distributed-таблицы, движки за пределами семейства MergeTree и цели с несколькими или резервными хостами внутри шарда, для которых нельзя доказать принадлежность одной логической таблице, отвергаются. Однохостовый нереплицируемый MergeTree без блочной дедупликации принимается с предупреждением, но exactly-once вырождается в at-least-once. Окно дедупликации конечно и вытесняется: если реплей доходит до ClickHouse позже вытеснения токена дедупликации, exactly-once также вырождается в at-least-once. Если известное окно дедупликации короче replay_horizon, синк пишет структурированное предупреждение с атрибутами Database, Table, DedupWindowSetting, ReplayHorizon и DedupWindow.
Сводная таблица гарантий по классам синка:
|
Exactly-once синк |
At-least-once синк |
At-most-once синк |
|
|
Промежуточное хранение в YTsaurus |
Да ( |
Нет |
Да при |
|
Дедупликация в ClickHouse |
|
Нет |
Нет |
|
Задержка |
Выше (доп. запись) |
Ниже |
Ниже |
|
Потеря при сбое |
Нет |
Нет |
Возможна после начала |
|
Дубликаты при сбое |
Нет |
Возможны |
Нет |
Подробнее — ClickHouse-расширение.
Сервис-лог
- Source: exactly-once — дедупликация по диапазонам чтения. Сервис-лог — это бесконечный источник, циклически обходящий таблицу. Гарантии exactly-once применяются в рамках каждого цикла обхода.
- Синк отсутствует. Сортированная динамическая таблица, используемая как источник сервис-лога, зачастую наполняется самим пайплайном через работу с внешним стейтом. Например, компьютейшен типа
TTransformComputationобновляет записи таблицы на каждой эпохе.
Подробнее — Сервис-лог.
Отказоустойчивость
- Пайплайн переживает выпадение отдельных машин и датацентров.
- Пока пайплайн остаётся активным и планирование продолжается, при сбое джобы контроллер удаляет её, отменяет Lease, если та ещё активна, и планирует необходимые новые джобы для затронутых текущих партиций, создавая для каждой новой джобы новую Lease. Подробнее о барьере и восстановлении см. в разделе Lease-транзакция.
- Внутренние потоки хранятся в динамических таблицах YTsaurus — потеря данных невозможна при штатной работе хранилища.
- Автоматическая балансировка партиций перераспределяет нагрузку при изменении топологии кластера.