Глоссарий YTsaurus Flow
- Введение
- Основная модель
- Timer (таймер)
- Resource (ресурс)
- Distributed Throttler (распределённый троттлер)
- State (стейт)
- Параллелизм и исполнение
- Внешние подключения
- Архитектура системы
- Время и корректность
- Обеспечение Exactly Once
- Дедупликация сообщений и порядок обработки сообщений
- Управление пайплайном
- См. также
Введение
В данной статье собраны все ключевые понятия YTsaurus Flow: от базовой модели обработки данных до архитектуры и механизмов обеспечения корректности. Материал выстроен от простого к сложному: сначала рассматривается потоковая обработка, а затем — параллелизм, внешние подключения, устройство системы и управление пайплайном.
При первом знакомстве с системой статью полезно прочитать от начала до конца, переходя по ссылкам за деталями. В дальнейшем страница послужит справочником: другие разделы документации часто ссылаются сюда, чтобы напомнить значение терминов.
Основная модель
Для понимания Flow как системы потоковой обработки данных можно провести аналогию с заводским конвейером:
- Пайплайн — целый завод, решающий конкретную бизнес-задачу.
- Стрим — лента конвейера, по которой движутся данные.
- Сообщение — отдельный элемент на ленте.
- Компьютейшен — рабочая станция, которая берёт элементы с одной или нескольких лент, обрабатывает их и кладёт результаты на другие ленты.
Pipeline (пайплайн)
Это конкретная выполняемая бизнес-задача. Может содержать множество различных узлов группировки и обработки данных.
С точки зрения YTsaurus пайплайн представлен Cypress-объектом типа pipeline. Подробности про устройство объекта, способы создания и управления — в разделе Объект Pipeline.
Stream (стрим)
Стрим — схематизированный именованный поток данных, соединяющий компьютейшены в графе пайплайна. Каждый стрим имеет фиксированную схему (набор типизированных полей) и идентифицируется по stream_id. Стрим состоит из сообщений. У каждого стрима ровно один писатель — компьютейшен, который выводит в него данные — и произвольное количество читателей.
Порядок сообщений в потоке в общем случае не гарантируется. Однако для производных сообщений с совпадающими ключами по всей цепочке lineage существуют гарантии относительного порядка. Подробнее — в разделе Порядок обработки сообщений.
Computation (компьютейшен)
Компьютейшен — узел графа пайплайна, выполняющий конкретное преобразование данных. Читает сообщения из входных стримов, обрабатывает их и записывает результаты в выходные стримы. Может иметь как множество входов, так и множество выходов. Входной поток разбивается на партиции для параллельной обработки.
Подробнее в разделе Computation.
Passthrough Computation (passthrough-компьютейшен)
Вид компьютейшена, не содержащий пользовательской бизнес-логики: входящие сообщения конвертируются в схему выходного стрима и передаются дальше без изменений. Используется для простого приведения схем между стримами.
Подробнее — Computation.
Swift
Принцип обработки данных во Flow, при котором результат работы компьютейшена не сохраняется в YTsaurus. Вместо этого функция преобразования должна быть строго детерминированной — при необходимости результат вычисляется повторно. Это снижает нагрузку на YTsaurus при сохранении гарантий exactly-once.
Подробнее в разделе Swift.
Message (сообщение)
Одно сообщение в рамках потока.
В контексте одного компьютейшена бывает несколько видов сообщений:
input— сообщения от других компьютейшенов, по сути,outputдругого компьютейшена, но со сгруппированной схемой.source— сообщения из внутренних сорсов.timer— специальные внутренние сообщения-таймеры.output— результат работы компьютейшена, который становится публичным и может попадать во внутренний синк и в другие компьютейшены.
Key (ключ группировки)
Набор значений полей сообщения, определяемых через group_by_schema в спеке компьютейшена. Все сообщения с одинаковым ключом направляются в одну партицию — каждой партиции назначен диапазон ключей [LowerKey; UpperKey). Ключ является единицей изоляции для стейта, таймеров и гарантий порядка.
Понятие ключа зависит от типа сообщения:
- Для сообщений типа
input— ключ определяется черезgroup_by_schema. - Для сообщений типа
source— ключ является партиционным ключом источника. - Для сообщений типа
output— ключа нет.
Lineage (родословная)
Совокупность входных сообщений и таймеров, из которых были получены конкретные выходные результаты компьютейшена. Фреймворк использует lineage для вычисления системных выходных сообщений и для обеспечения гарантий порядка производных сообщений.
Подробнее в разделе Lineage.
Timer (таймер)
Механизм отложенного вызова, привязанный к конкретному ключу группировки в рамках TTransformComputation. Таймер позволяет компьютейшену сказать: «разбуди меня, когда время X наступит».
Каждый таймер содержит два временных поля:
TriggerTimestamp— момент срабатывания. Когда EventWatermark (или другой сконфигурированный вотермарк) превысит это значение, таймер передаётся в компьютейшен для обработки.EventTimestamp— бизнес-время исходного события, «придержанное» в таймере.
Таймеры надёжно сохраняются в YTsaurus. Результат обработки таймера фиксируется в той же транзакции, в которой таймер удаляется, — это обеспечивает гарантии exactly-once.
Важно
В текущей реализации все активные таймеры дополнительно хранятся в памяти процесса. При большом количестве таймеров это может приводить к ошибкам Out of Memory при старте джобов.
Подробнее в разделе Таймеры.
Resource (ресурс)
Ресурсы существуют для описания данных, общих для нескольких джобов. Это могут быть клиенты к YTsaurus и другим системам, кеши для обращений во внешние системы, модели машинного обучения и т. п.
На данный момент ресурсы могут быть только статическими (то есть неизменяемыми без перезапуска).
Distributed Throttler (распределённый троттлер)
Во Flow есть возможность создавать именованные распределённые троттлеры — общие для всех джобов пайплайна. Их используют для ограничения нагрузки на внешние API, выравнивания пропускной способности между партициями одного компьютейшена или притормаживания чтения из источников. Реализованы как token bucket на контроллере; джобы запрашивают квоту перед обработкой сообщений или вручную из пользовательского кода. Подробнее в разделе Distributed Throttler.
State (стейт)
Персистентные данные, связанные с конкретным ключом группировки в рамках компьютейшена. Хранятся в динамических таблицах YTsaurus и обновляются атомарно при коммите эпохи. Пустое значение стейта соответствует отсутствию строки в таблице.
Подробнее в разделе Stateful processing.
Параллелизм и исполнение
Для обеспечения масштабируемости процесса обработки данных, Flow разбивает каждый компьютейшен на множество партиций, которые обрабатываются параллельно в рамках джобов.
Partition (партиция)
Входной поток в конкретный компьютейшен может быть достаточно велик, поэтому для целей параллельной обработки данных компьютейшен разбивается на множество партиций. Ориентировочные ограничения по потоку для одной партиции: не более 1 MB/s и не более 1000 сообщений в секунду.
За счёт поля group_by_schema в спеке компьютейшена входной поток группируется и данные по одному ключу всегда будут попадать в конкретную партицию. Для этого за каждой партицией определён диапазон ключей [LowerKey; UpperKey).
Job (джоб)
Конкретный запуск обработчика конкретной партиции на воркере, выбранном контроллером.
Epoch (эпоха)
Джоб обрабатывает данные эпохами для регулярного надёжного сохранения прогресса обработки данных. Механизм эпох позволяют поддерживать гарантии exactly-once при разумной нагрузке на YTsaurus, а также для минимизации времени простоя при восстановлении от сбоев.
Некоторые эпохи могут не содержать коммитов в YTsaurus. Некоторые реализации компьютейшена могут вообще не делать коммиты.
Layout
Список всех джобов и партиций в системе. Иными словами — описание, что где сейчас работает.
Внешние подключения
Flow взаимодействует с внешними системами (очередями, таблицами и т. д.) через коннекторы. Каждый коннектор предоставляет сорс для чтения и/или синк для записи.
Connector (коннектор)
Компонент, обеспечивающий связь пайплайна с некоторой внешней системой через чтение сообщений из source коннектора и/или запись сообщений в sink коннектора. Коннектор может включать в себя и source, и sink (например, в случае QYT) или же только source (например, в случае servicelog).
В коде и в спеках описываются напрямую source'ы и sink'и. Сам коннектор не существует как сущность, это скорее элемент логической группировки source'ов и sink'ов по внешней системе, с которой идёт взаимодействие.
Подробнее о коннекторах можно прочитать в документации по коннекторам.
Source (сорс)
Компонент, позволяющий компьютейшену читать сообщения из внешней системы. Реализуется в виде отдельного класса и конфигурируется в спеке компьютейшена.
Sink (синк)
Компонент, позволяющий компьютейшену записывать сообщения во внешнюю систему. Реализуется в виде отдельного класса и конфигурируется в спеке компьютейшена.
Архитектура системы
Запущенный пайплайн состоит из набора процессов двух ролей — контроллера и воркера. Контроллер управляет распределением работы, а воркеры непосредственно обрабатывают данные.
Controller (контроллер)
Управляет всем пайплайном. Назначает джобы воркерам, предоставляет API для управления пайплайном, отвечает за синхронизацию всех настроек, следит за состоянием всех воркеров и собирает от них статусы. Может быть запущен в нескольких экземплярах (рекомендуемое — 2-3), тогда один из инстансов будет лидером, а остальные будут ожидать в режиме stand-by для обеспечения отказоустойчивости.
Спеки для выполнения задаются в контроллер через специальный API.
Computation Controller
Контроллер конкретного компьютейшена: отвечает за создание необходимого количества партиций с нужными настройками. Не путать с просто контроллером. Computation Controller — это класс, а контроллер — это программа, которая управляет работой кластера и создаёт необходимые Computation Controller-ы.
JobManager
Модуль, отвечающий за создание и распределение джобов по воркерам.
LeaseManager
Модуль, отвечающий за создание и пинг Lease — мастер-транзакций, которые используются джобами при коммитах в качестве prerequisite_transaction_ids. Позволяют обеспечить гарантии exactly-once.
Worker (воркер)
Один из двух основных видов процессов в системе, осуществляющий сами вычисления в рамках назначенных ему джобов. Регулярно шлёт хартбиты в контроллер, отправляя ему статусы всех джобов. В ответ получает актуальные настройки системы, в частности Layout.
Хартбит (heartbeat)
Периодическое сообщение, которое воркер отправляет контроллеру со статусами всех своих джобов. В ответе контроллера приходят актуальные настройки системы, в частности Layout; воркер, переставший слать хартбиты, считается недоступным.
Message Distributor
Модуль воркера, отвечающий за рассылку сообщений между различными джобами. Выполняется на базе текущего лэйаута системы. Шлёт сообщение до тех пор, пока не получит от получателя подтверждения обработки, включая факт надёжного коммита результата обработки (в кодовой базе упоминается как MarkPersisted).
Resource Manager
Компонент воркера, создающий ресурсы для общего пользования. Например, несколько компьютейшенов могут использовать одну и ту же бинарную базу или один и тот же клиент YTsaurus.
BufferStateManager
Компонент, отвечающий за создание и поддержание актуальных размеров всех (или почти всех) буферов в системе. У каждой джобы есть буферы входящих и исходящих сообщений. Задача BufferStateManager — обеспечить необходимый для бесперебойной работы размер буферов.
Companion (компаньон)
Отдельный процесс, запускаемый на одном хосте с воркером и выполняющий пользовательский код на Python, Go, Java или Kotlin либо C++. Взаимодействует с воркером по gRPC. Выполняет компьютейшены вне процесса воркера: как на языках, отличных от C++, так и с помощью C++-компаньона.
Подробнее в разделе Companion.
Время и корректность
Корректная обработка событий во времени — одна из ключевых задач потоковой системы. Flow отслеживает три вида временных меток у каждого сообщения и поддерживает механизм вотермарков для определения прогресса обработки.
EventTimestamp, SystemTimestamp, AlignmentTimestamp, StabilizedEventTimestamp и Watermarks
В отношении каждого сообщения во Flow могут использоваться следующие таймстемпы:
- EventTimestamp — истинное время события, ассоциированного с данным сообщением.
- SystemTimestamp — время создания конкретного сообщения. Для сообщений, сгенерированных внутри Flow, в качестве
SystemTimestampберётся «время YTsaurus». Дляsourceпотоков — время появления сообщения в удалённой системе. - AlignmentTimestamp — таймстемп, используемый для выравнивания прогресса обработки партиций. Вычисляется автоматически.
- StabilizedEventTimestamp — таймстемп, вычисляющийся на основе
AlignmentTimestamp, используется как неубывающее на потоке сообщений приближениеEventTimestamp(при совпадающем наборе ключей в lineage).
Первые три таймстемпа хранятся в самом сообщении, а StabilizedEventTimestamp вычисляется на лету при необходимости.
Вотермарк (watermark) — это временная метка, означающая, что система больше не ожидает поступления более старых событий. На каждом стриме в системе поддерживаются текущие SystemWatermark и EventWatermark.
Примечание
К сожалению, обычные протоколы записи данных в очереди не могут гарантировать, что в очередь не будет записано вчерашнее событие. Поэтому про SystemWatermark по внутренним потокам можно говорить, что он является абсолютным, а про EventWatermark — лишь то, что он эвристичен и является приближённым.
Обеспечение Exactly Once
Flow по умолчанию обеспечивает exactly-once семантику обработки событий. Механизм основан на барьерных Lease-транзакциях, дедупликации входных сообщений по message_id (для внутренних потоков) и по оффсетам (для источников), а также на атомарности выходных данных через транзакции эпохи. При необходимости семантику можно ослабить до at-least-once или at-most-once.
Подробнее — в разделе Гарантии обработки.
Дедупликация сообщений и порядок обработки сообщений
Основным способом дедупликации сообщений является дедупликация по message_id, а не по оффсету. Механизм с оффсетами используется только при работе с очередями. В том числе, Flow не даёт никакой гарантии на порядок обработки событий.
Управление пайплайном
Этот раздел описывает операции жизненного цикла пайплайна: от описания его конфигурации через спеки до релиза, миграции и работы в различных окружениях.
Spec и DynamicSpec
Любой выполняемый на Flow пайплайн представляет собой комбинацию из двух YSON конфигов, называемых Spec и DynamicSpec.
Spec— определяет статические свойства пайплайна, его топологию, свойства узлов, связи между ними, типы объектов и т. п. Статическая спека может меняться только при условии остановки пайплайна.DynamicSpec— определяет динамические свойства пайплайна.
Соответственно, спеки содержат свойства всех компонент системы, как пользовательского слоя (например, список компьютейшенов и схемы стримов между ними), так и настроек системного уровня (кодеки сжатия в системных таблицах, размеры буферов и т. п.).
Релиз пайплайна
Обновление пайплайна, которое включает изменение его исполняемых файлов.
Обновление спек пайплайна
Обновление статической и/или динамической спеки пайплайна.
Запуск и остановка пайплайна
Изменение целевого состояния пайплайна и переход пайплайна в это состояние. Существующие состояния:
unknown— пайплайн ещё не запущен.working— пайплайн работает, сообщения обрабатываются.stopped— пайплайн остановлен, все сообщения дообработаны.paused— пайплайн приостановлен (джобы остановлены, промежуточные сообщения могут быть не дообработаны).draining— переходное состояние: пайплайн в процессе остановки.pausing— переходное состояние: пайплайн в процессе приостановки.completed— финальное состояние: все источники пайплайна были конечными (finite = true) и все сообщения из них обработаны. Выйти из этого состояния нельзя — потребуется пересоздать пайплайн. Чаще всего встречается в интеграционных тестах, а также может возникать в продакшн-пайплайнах при некорректной последовательности действий во время деплоя (например, если источник был ошибочно помечен как конечный) или при выставлении некорректной спеки.
Внутренние таблицы пайплайна
Служебные таблицы, которые размещаются в директории пайплайна в YTsaurus. Они необходимы для работы пайплайна (например, для дедупликации входных сообщений, хранения информации по партициям и т. д.).
Полный список внутренних таблиц и их назначение приведены в разделе Объект Pipeline → Внутренние таблицы.
Пользовательские таблицы
Таблицы, с которыми работают пользователи. Как правило, их содержимое имеет продуктовый смысл, например:
- Упорядоченные таблицы — очереди, из которых пользовательский пайплайн читает входные события или в которые пишет выходные.
- Сортированные таблицы с состояниями сущностей (профили баннера/оффера/пользователя и т. д.), которые пользователь строит в процессе работы пайплайна.
Миграция
Процесс изменения структуры объектов в YTsaurus — например, добавление новых таблиц, изменение схемы существующих таблиц, их удаление.
Основные причины для выполнения миграций:
- Выкатка новой версии пайплайна, не совместимой с текущим состоянием объектов в YTsaurus. Миграция позволяет без потери данных обновить объекты до состояния, с которым сможет работать новая версия пайплайна.
- Изменение типа таблицы (replicated/standalone/chaos) из соображений производительности или доступности.
Примеры изменений:
|
Требуют миграции
|
Не требуют миграции
|
Окружение (environment)
Инсталляция, где запускается пайплайн, — тестинг, препродакшн или продакшн.
Является синонимом понятия «стейдж» (stage).
См. также
Можно менять только при условии остановки пайплайна.