Swift в YTsaurus Flow

Swift — принцип обработки данных в YTsaurus Flow, при котором результат работы компьютейшена не сохраняется в YTsaurus. Вместо этого функция преобразования должна быть строго детерминированной: если результат потребуется повторно (например, при перезапуске джоба), он будет вычислен заново из тех же входных данных.

Зачем это нужно

В классическом подходе (TTransformComputation) каждая эпоха завершается транзакционной записью результатов в YTsaurus. Это обеспечивает exactly-once, но создаёт нагрузку на кластер: для каждого входного сообщения выполняется лукап и запись в таблицу дедупликации, плюс по одной записи на каждое выходное сообщение.

Swift снимает это ограничение: если функция детерминирована, хранить её вывод не нужно — при необходимости его можно воспроизвести. Это позволяет:

  • снизить нагрузку на YTsaurus до нуля или до минимума (только метаданные),
  • повысить пропускную способность при stateless-преобразованиях.

Как сохраняются гарантии exactly-once

Несмотря на отсутствие записи выходных данных в YTsaurus, гарантии exactly-once сохраняются за счёт детерминированности:

  • Если джоб упал до доставки результата, Flow перезапускает его и получает тот же вывод из тех же входных данных.
  • Message Distributor продолжает доставлять сообщение до получения подтверждения (MarkPersisted) от получателя, что исключает потери.

Таким образом, exactly-once обеспечивается не хранением вывода, а идемпотентностью вычисления.

Требование детерминированности

Функция преобразования в Swift-компьютейшене должна быть детерминированной: при одних и тех же входных данных должен возвращаться одинаковый вывод, включая порядок сообщений.

Требование распространяется и на lineage — привязку выходных сообщений к родительским: при повторном вычислении каждый вывод должен получить тех же родителей в том же порядке. Типичная ошибка — итерация по неупорядоченной структуре (например, hash-таблице или множеству) при разбиении батча на группы: порядок групп меняется от запуска к запуску, и повторное вычисление даёт другой вывод. По умолчанию у каждого выходного сообщения Swift-компьютейшена ровно один родитель; несколько родителей допустимы только при allow_batching_with_relaxed_guarantees.

Важно

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

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

По умолчанию Flow дописывает к message id производного сообщения его порядковый номер. Если два повтора Swift-вычисления порождают на одной позиции разные сообщения, они получают одинаковый message id. Это может привести к потере данных или появлению дубликатов, а если у сообщений различаются ключи — также к нарушению внутренних инвариантов Flow. В C++, Java, Python и Go через output options можно исключить порядок из идентичности сообщения, выбрав хеш payload или пользовательский суффикс (см. описание C++ API). Это защищает от такого совпадения идентификаторов только при стабильном семантическом суффиксе и не делает недетерминированное вычисление детерминированным.

Классы Swift-компьютейшенов

В Flow реализовано два базовых Swift-класса:

TSwiftMapComputation

Детерминированный Map без материализации результатов в YTsaurus.

  • Нагрузка на YTsaurus: ~0 записей за эпоху (только системные фоновые процессы).
  • Не поддерживает: Source, Sink.
  • Поддерживает: таймеры и key-visitor-стримы — только для работы со стейтом, например для фонового клинапа (GC). Эмит выходных сообщений из обработки таймера или визита запрещён: output-стрим не может зависеть от таймер- или visit-стрима в streams_dependency. Так как таймер-стримы по умолчанию добавляются в зависимости каждого output-а, спека с таймерами и output-ами обязана задавать streams_dependency явно.
  • Требует: строгой детерминированности и того, чтобы у каждого результирующего сообщения был ровно один родитель — входное сообщение.

Подробнее — в разделе Computation (C++).

TSwiftPassthroughComputation

Passthrough-компьютейшен — наследник TSwiftMapComputation. Конвертирует input-сообщения в схему output-стрима без пользовательской логики. Подробнее — Computation (C++).

TSwiftOrderedSourceComputation

Основной класс для чтения упорядоченных данных из внешних источников.

  • Нагрузка на YTsaurus: ~1–2 записи за эпоху (метаданные для восстановления; сами сообщения не сохраняются).
  • Поддерживает: WatermarkStrategy для оценки вотермарков.
  • Требует: ровно одного Source, реализующего IOrderedSource.

Подробнее — в разделе Computation (C++).

TSwiftPassthroughOrderedSourceComputation

Passthrough-компьютейшен — наследник TSwiftOrderedSourceComputation. Преобразует source-сообщения в output-стрим приведением к новой схеме. Подробнее — Computation (C++).

Сравнение с TTransformComputation

Тип Запись в YTsaurus за эпоху Поддержка таймеров Поддержка стейта Требование детерминированности
TTransformComputation 2 на каждое входное сообщение и 1 на выходное Да Да Нет
TSwiftOrderedSourceComputation ~1–2 (метаданные) Нет Нет Да
TSwiftPassthroughOrderedSourceComputation ~1–2 (метаданные) Нет Нет Да
TSwiftMapComputation ~0 Да (только стейт) Да Да
TSwiftPassthroughComputation ~0 Нет Нет Да

См. также