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 | Нет | Нет | Да |