Статические таблицы в YTsaurus Flow
- NYT::NFlow::TUnitedParameters<NYT::NFlow::NStaticTableConnector::TSource>
- NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnector::TSource>
- NYT::NFlow::TUnitedParameters<NYT::NFlow::NStaticTableConnector::TArrivalOrderTableSink>
- NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnector::TArrivalOrderTableSink>
- Импорт статической таблицы в стейт
- Альтернативы
- См. также
Коннектор к статическим таблицам YTsaurus.
Код коннектора находится здесь.
Статические таблицы — особый вид источника. В нём изначально нет партиций, у строк в них нет таймстемпов записи, сами таблицы как правило иммутабельны и в них не происходит дозаписи. При этом часто читать нужно бесконечную последовательность таблиц. Так же "запись" в этот источник идёт очень крупными блоками, поэтому читать из него необходимо с ограничением скорости, чтобы не отнимать ресурсы у более важных сорсов пайплайна.
В связи с этим основная сложность этого источника лежит в контроллере, которому необходимо понять, какие таблицы читать, какие таймстемпы (SystemTimestamp, EventTimestamp) для них выбрать, что делать, если таблица неожиданно исчезает и так далее.
Плановые timestamps чтения
Установите use_planned_timestamps = %true в статических параметрах источника, чтобы SystemTimestamp и EventTimestamp равнялись плановому времени начала чтения диапазона строк. Timestamp пропорционален доле строк, уже выданных на чтение. Контроллер фиксирует начало и общую длительность чтения при запуске таблицы по текущим лимитам. Рестарты процессов сохраняют план. При смене реплики, возврате на неё или повторном чтении таблицы план оставшихся диапазонов пересчитывается от текущего времени по текущим лимитам; порядок перечитывания сохраняется. Timestamps могут оказаться в будущем.
Минимальные timestamps учитывают незавершённые диапазоны, включая ещё не распределённые. Опция выключена по умолчанию. Включение и выключение действует со следующей таблицы. После выключения следующая таблица снова использует timestamps локаторов. Сообщения с timestamps ниже опубликованного watermark считаются опоздавшими; это ожидаемое поведение.
Настройки сорса
Класс сорса: NYT::NFlow::NStaticTableConnector::TSource.
Статическая спека:
NYT::NFlow::TUnitedParameters<NYT::NFlow::NStaticTableConnector::TSource>
Источник: yt/yt/flow/library/cpp/common/registry-inl.h
|
Параметр |
Описание |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: TDuration |
Дополнительные параметры
|
|
Тип: TDuration |
|
|
Тип: |
Динамическая спека:
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnector::TSource>
Источник: yt/yt/flow/library/cpp/common/registry-inl.h
|
Параметр |
Описание |
|
|
Тип: TDuration |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: |
Дополнительные параметры
Эти параметры для тонкой настройки, не рекомендуется трогать без глубокого понимания системы.
|
|
Тип: TDuration |
|
|
Тип: |
|
|
Тип: NYT::NYTree::TSize |
|
|
Тип: TDuration |
|
|
Тип: TDuration |
|
|
Тип: |
|
|
Тип: |
|
|
Тип: TDuration |
Запись в статические таблицы в порядке прихода
Класс синка: NYT::NFlow::NStaticTableConnector::TArrivalOrderTableSink.
Синк создаёт непрерывную последовательность таблиц с фиксированным шагом table_period. Непустой текущий слот закрывается по границе времени либо при достижении max_row_count/max_data_weight; закрытие по лимиту также сдвигает следующий логический таймстемп на один период, поэтому последовательность может обогнать wall clock. Пока последовательность впереди wall clock, закрытие по времени не срабатывает: неполный батч будет записан, только когда wall clock догонит текущий логический таймстемп, а рестарт этого не сбрасывает — отставание хранится во внешнем стейте. Величина задержки пропорциональна всплеску: число слотов, закрытых по лимиту, умноженное на table_period. Пустой слот T создаётся строго по порядку, только если T <= wall clock и при этом известный ненулевой системный watermark входного stream строго больше T + table_period. Поэтому равенства watermark границе недостаточно, а пустые таблицы в будущем не создаются.
Таблица и прогресс коммитятся одной master-транзакцией. Прогресс лежит в атрибуте @progress самой output_directory и содержит владельца (pipeline, computation, sink id), общую последовательность таблиц и отдельный frontier (system_timestamp, message_id) для каждой партиции. Все партиции делят одну последовательность таблиц, а frontier используется для дедупликации при replay: при частично покрытом replay синк записывает только непокрытый хвост без перезапуска job. Callback доставки вызывается только после успешного внешнего коммита и следующего коммита Flow. Писатели прогресса разводятся shared-локом с ключом атрибута progress, поэтому создание выходных таблиц в той же директории не блокируется.
Каждому синку нужна собственная output_directory: если в атрибуте записан другой владелец, синк падает с ошибкой и просит удалить атрибут вручную. Перед передачей директории другому пайплайну остановите пишущий пайплайн; новый владелец продолжит сетку после самой свежей таблицы в директории. Директорию с детьми без атрибута table_timestamp синк не принимает: он падает с ошибкой, пока их не уберут. В частности, main и DLQ одного reader должны писать в разные директории. Frontier партиции удаляется, как только её system_timestamp опускается ниже системного watermark входного stream — к этому моменту партиция гарантированно доставила всё, что произвела, поэтому её frontier больше не нужен.
output_directory может находиться на кластере, отличном от кластера пайплайна. Обязательный параметр table_ttl задаёт время жизни выходных таблиц. Для каждой таблицы Cypress expiration_time равен max(table_timestamp, время создания таблицы) + table_ttl: при догоняющей обработке таблица не удалится сразу, даже если её table_timestamp давно прошёл. Когда последовательность таблиц соответствует реальному времени, в директории остаётся не более table_ttl / table_period таблиц; при создании множества догоняющих таблиц за короткое время их может временно быть больше. Спека с отношением table_ttl / table_period больше 40000 отклоняется из-за лимита Cypress на число детей.
Синк требует ровно одного входного stream и возрастающих MessageId внутри него, поэтому применим в Transform- и SwiftOrderedSource-режимах. Партиция идентифицируется по SourceKey, если он есть, иначе по PartitionId. Схема выходных таблиц берётся из входного stream. Вес сообщения для лимита max_data_weight по умолчанию равен размеру сообщения; опциональная data_weight_column (тип int64 или uint64, неотрицательные значения) задаёт пользовательский вес, а null в ней означает размер сообщения. Транзакция инициализации и коммита повторяется до успеха или отмены job. Раннее создание пустых таблиц без входных сообщений требует API инициализации синков, которого пока нет у process function; если это обязательное требование, оформите отдельный запрос на расширение API, а не создавайте raw Computation.
Статическая спека:
NYT::NFlow::TUnitedParameters<NYT::NFlow::NStaticTableConnector::TArrivalOrderTableSink>
Источник: yt/yt/flow/library/cpp/common/registry-inl.h
|
Параметр |
Описание |
|
|
Тип: |
|
|
Тип: TDuration |
|
|
Тип: TDuration |
|
|
Тип: |
|
|
Тип: |
Динамическая спека:
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnector::TArrivalOrderTableSink>
Источник: yt/yt/flow/library/cpp/common/registry-inl.h
|
Параметр |
Описание |
|
|
Тип: |
|
|
Тип: |
Дополнительные параметры
Импорт статической таблицы в стейт
Типичный сценарий: батч-процесс (YQL, …) периодически пересобирает справочник в виде статической YT-таблицы, а реалтайм-пайплайн должен джойнить свой поток против свежей версии этого справочника. Коннектор static_table поддерживает этот сценарий "из коробки", позволяя загрузить такую таблицу в стейт пайплайна с произвольной бизнес-логикой над каждой строкой перед записью.
Как это работает
Схема: один computation читает справочник через static_table-сорс и пишет его в стейт (loader); второй computation на входном потоке достаёт значение по ключу из того же стейта и эмитит обогащённое сообщение (enricher).
static_table-сорс потоково отдаёт строки таблицы как обычные сообщения. Контроллер сам определяет, какая версия таблицы свежая, и подаёт её строки в пайплайн в виде отдельного стрима.- Loader-computation (stateful) принимает каждую строку, при необходимости валидирует/нормализует её (
trim,lower-case, обогащение из других источников, фильтрация и т. п.) и пишет результат в стейт по ключу. - Enricher-computation, висящий на потоке событий, читает значения из того же стейта по ключу события через read-only joiner и эмитит обогащённое сообщение в синк.
Главное преимущество этого подхода — возможность написать бизнес-логику обработки каждой строки после сборки таблицы, но до того, как она попала в стейт. Это удобно, когда форма данных на диске не совпадает с тем, что нужно джойну, или когда часть строк надо отфильтровать/обогатить из других сорсов на лету.
Семантика перезаливки
Каждая новая версия статической таблицы — полная перезаливка всего справочника в стейт. Автоматической очистки удалённых строк у этой схемы нет: если в новой версии строка пропала, её нужно явным образом удалить из стейта — например, помечая строки версией и периодически удаляя устаревшие через механизм key visit.
Свежая версия выбирается контроллером по правилам, заданным в TTableTimestampLocatorSpec (см. статическую спеку сорса выше) — чаще всего это самая свежая таблица в директории, имя которой парсится как ISO8601-таймстемп. Появление новой таблицы в директории автоматически инициирует новый прогон.
Когда выбирать
- Над каждой строкой нужна бизнес-логика (нормализация, валидация, обогащение).
- Допустима полная перезаливка на каждый ребилд.
Эти условия не жёсткие: подход работает и без бизнес-логики над строками, а сложности полной перезаливки часто обходятся (см. про key visit выше). Тем не менее стоит взвесить и другие варианты — см. раздел Альтернативы.
Конфигурация (набросок)
Спека целиком — в примерах ниже; здесь только ключевая обвязка пайплайна. Loader пишет в стейт через manager (read-write), enricher читает его через joiner (read-only).
Набросок спеки
{
"spec" = {
"computations" = {
"reference_reader" = {
"computation_class_name" = "...";
"output_stream_ids" = ["reference"];
"source_streams" = {
"reference_table" = {
"source_class_name" = "NYT::NFlow::NStaticTableConnector::TSource";
"parameters" = { "tables_path" = "<cluster=primary>//path/to/reference"; };
};
};
};
"reference_loader" = {
"computation_class_name" = "...";
"input_stream_ids" = ["reference"];
"group_by_schema" = [
{"name" = "hash"; "expression" = "farm_hash(key)"; "type" = "uint64"; required = %true;};
{"name" = "key"; "type" = "uint64";};
];
"external_state_managers" = {
"/reference_state" = {
"external_state_manager_class_name" = "NYT::NFlow::TSimpleExternalStateManager";
"parameters" = { "path" = "<cluster=primary>//path/to/state"; };
};
};
};
"enricher" = {
"computation_class_name" = "...";
"input_stream_ids" = ["event"];
"group_by_schema" = [
{"name" = "hash"; "expression" = "farm_hash(key)"; "type" = "uint64"; required = %true;};
{"name" = "key"; "type" = "uint64";};
];
"external_state_joiners" = {
"/reference_state" = {
"external_state_joiner_class_name" = "NYT::NFlow::TSimpleExternalStateJoiner";
"parameters" = { "path" = "<cluster=primary>//path/to/state"; };
};
};
};
};
};
}
Полные рабочие конфиги, бинарь, yt_sync и интеграционные тесты — в примерах.
Примеры
Полностью рабочие пайплайны с тестами, демонстрирующие данный подход:
- C++:
examples/cpp/static_table_join - Python:
examples/python/static_table_join - Java:
examples/java/static_table_join - Kotlin:
examples/kotlin/static_table_join - Go:
examples/go/static_table_join
Альтернативы
static_table расширение → стейт — не единственный способ организовать джойн со справочником; альтернативы и их компромиссы — в таблице ниже.
|
Подход |
Когда выбирать |
Цена |
|
нужна бизнес-логика над строкой до записи в стейт; допустима полная перезаливка |
стоимость перезаливки внутри пайплайна на каждой версии |
|
|
#2 Конвертация в динамическую таблицу + symlink под external state |
большой объём, lookup по запросу, нужна атомарная смена версии данных |
сетевой lookup в динтаблицу + эксплуатация симлинка |
|
ноль сетевых обращений на джойне, объём ограничен памятью/диском воркера |
сложная эксплуатация, формат и доставка БД на воркеры |
Конвертация в динамическую таблицу + symlink под external state
В этом подходе роль стейта играет сама динамическая таблица: батч-процесс собирает новую версию справочника как сортированную динамическую таблицу …/reference.vN, монтирует её, а пайплайн смотрит на неё через Cypress-симлинк …/current → …/reference.vN. После сборки следующей версии симлинк атомарно перенаправляется на …/reference.v(N+1) (yt link --force …/reference.v(N+1) …/current либо set @target_path), и пайплайн начинает получать новые значения.
В пайплайне path коннектора external-state ссылается на симлинк, а не на конкретную версию. Лукап выполняется read-only joiner-ом (TSimpleExternalStateJoiner); по умолчанию joiner перечитывает значение из YT на каждый лукап, поэтому переключение симлинка сразу видно пайплайну без рестарта. Подробнее про подготовку динамической таблицы и её схему — в sorted-dynamic-table.md.
Примеры
- C++:
examples/cpp/external_state_join - Python:
examples/python/external_state_join - Java:
examples/java/external_state_join - Kotlin:
examples/kotlin/external_state_join - Go:
examples/go/external_state_join
В тестах примеров yt_sync сначала собирает reference.v1 и линкует current → reference.v1, пайплайн обогащает событие значением v1, затем собирается reference.v2, симлинк атомарно перенаправляется на новую версию, и для того же ключа пайплайн начинает отдавать значение v2.
Встроенная БД и деплой через Resource
Если справочник целиком помещается в память или на локальный SSD воркера, и доступ к нему нужен совсем без сетевых походов, имеет смысл собирать его как встроенную БД: батч-процесс выкатывает готовый файл/директорию, а пайплайн поднимает её как локальный лукап-движок прямо внутри процесса джоба.
Достоинства — нулевые сетевые обращения на джойне, минимальная задержка лукапа (только локальный диск/память), полная независимость от внешних KV-сервисов. Недостатки — объём ограничен ресурсами воркера; формат БД, её схема и совместимость версий ложатся на команду; доставку новой версии (Resource, переподготовка, чек целостности) нужно эксплуатировать самостоятельно. Подходит для относительно небольших справочников (единицы гигабайт), которые обновляются нечасто.