Статические таблицы в YTsaurus Flow

Коннектор к статическим таблицам 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

Параметр

Описание

finite

Тип: bool
Значение по умолчанию: false
Считать source конечным, то есть запомнить число сообщений в нём при старте и перевести stream в состояние completed при вычитывании этого числа сообщений.

tables

Тип: std::optional<std::vector<NYT::NYPath::TRichYPath>>
Список таблиц для чтения с указанием кластеров. Этот параметр альтернативен tables_path. Каждый элемент должен быть таблицей: симлинк не разыменовывается до целевой таблицы, а любой нетабличный узел приводит к ошибке.

tables_path

Тип: std::optional<NYT::NYPath::TRichYPath>
Путь до директории со статическими таблицами с указанием одного кластера или упорядоченного списка кластеров-реплик. Этот параметр альтернативен tables.

table_name_filter

Тип: NYT::TIntrusivePtr<NYT::NRe2::TRe2>
Регулярное выражение для фильтрации таблиц по имени: читаются только те таблицы, имя которых совпадает с выражением. Синтаксис выражений — RE2.

event_timestamp_locator

Тип: NYT::TIntrusivePtr<NYT::NFlow::NStaticTableConnector::TTableTimestampLocatorSpec>
Значение по умолчанию: {'attribute': 'key'}
По умолчанию берёт таймстемп из имени таблицы. Это время соответствует времени создания данных, оно будет проброшено в EventTimestamp сообщений. Таблицы с одинаковым таймстемпом упорядочиваются персистентно; новые таблицы не должны появляться позади уже обработанного event-time frontier.

system_timestamp_locator

Тип: NYT::TIntrusivePtr<NYT::NFlow::NStaticTableConnector::TTableTimestampLocatorSpec>
Значение по умолчанию: {'attribute': 'creation_time'}
По умолчанию берёт таймстемп из времени создания таблицы. Это время соответствует времени записи данных в сорс. То есть моменту, когда пайплайн может увидеть эти данные и начать читать. Это время пробрасывается в SystemTimestamp сообщений.

use_planned_timestamps

Тип: bool
Значение по умолчанию: false
Использовать сохранённое плановое время начала чтения диапазона для обоих timestamps входных сообщений. План фиксируется при запуске таблицы контроллером и сохраняется при рестартах и изменении лимитов скорости. Локаторы timestamps по-прежнему определяют обнаружение и порядок таблиц. По умолчанию режим выключен.

ignore_symlinks

Тип: bool
Значение по умолчанию: false
Флаг, который позволяет игнорировать симлинки внутри папки с таблицами.

skip_non_table_nodes

Тип: bool
Значение по умолчанию: false
Пропускать нетабличные узлы во входной директории вместо ошибки.

idle_watermark_delay

Тип: std::optional<TDuration>
Значение по умолчанию: 3600000
Задержка продвижения ватермарка от текущего времени, когда source не читает таблицу. Ватермарк читаемой таблицы сразу продвигается до её event timestamp без вычитания задержки. Значение # отключает продвижение от текущего времени; новые таблицы при этом продолжают продвигать ватермарк.

failover_delay

Тип: TDuration
Значение по умолчанию: 5m
Для реплицированного входа — время недоступности активного кластера до переключения текущей таблицы.

Дополнительные параметры

update_info_period

Тип: TDuration
Значение по умолчанию: 15s
Период обновления служебной информации source о партиции. Конкретный коннектор может использовать этот тик для запросов статуса и проверки живости сессии.

byte_size_alpha

Тип: double
Значение по умолчанию: 0.05
Коэффициент экспоненциального сглаживания средней суммы байтов и числа сообщений на один offset: чем он больше, тем быстрее оценка реагирует на новые данные.

Динамическая спека:

NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnector::TSource>

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

desired_table_process_time

Тип: TDuration
Значение по умолчанию: 1h
За какое время читать одну входную таблицу. Основной параметр для контроля скорости чтения.

max_rows_per_second

Тип: double
Значение по умолчанию: 1000000000000.0

max_bytes_per_second

Тип: double
Значение по умолчанию: 1000000000000.0

max_event_timestamp

Тип: std::optional<unsigned long>

restart_instant

Тип: TInstant
Значение по умолчанию: 1970-01-01T00:00:00.000000Z
Установите в этот параметр текущее время в iso8601, чтобы забыть текущий прогресс и начать читать статические таблицы заново. Значение параметра можно безопасно уменьшать, это не приведёт к дополнительному рестарту чтения. Рестарт происходит, только если текущий restart_instant больше, чем последний restart_instant, что сорс сохранил внутри себя. При этом параметр никак не соотносится с таймстемпами статических таблиц и не производит их фильтрацию для рестарта чтения.

allow_v1_migration

Тип: bool
Значение по умолчанию: true
Управляет выходом из V1-совместимого порядка. Значение по умолчанию — true. Checkpoint формы V1 распознаётся как V1, после чего source с включённым флагом необратимо переводит его через V1 → Draining → V2 или сразу в V2, если V1-таблица не обрабатывается. Значение false откладывает первоначальный переход. Персистентные Draining и V2 остаются авторитетными и не понижаются при смене флага на false.

Дополнительные параметры

Эти параметры для тонкой настройки, не рекомендуется трогать без глубокого понимания системы.

unavailable_threshold

Тип: TDuration
Значение по умолчанию: 5m
Сколько времени подряд источник должен быть недоступен, чтобы партиция считалась стабильно недоступной. Засчитывается только то время, когда джоб работал и видел ошибку: простой между перезапусками в него не попадает, а любой успешный ответ источника обнуляет накопленное.

min_event_timestamp

Тип: std::optional<unsigned long>
Таблицы у которых EventTimestamp меньше MinEventTimestamp, не будут процесситься.

max_partition_count

Тип: NYT::NYTree::TSize
Значение по умолчанию: 10K
Менять не рекомендуется. Ограничение сверху на число одновременно живущих партиций.

throttler_period

Тип: TDuration
Значение по умолчанию: 10s
Менять не рекомендуется. Регулирует окно на котором регулируется скорость чтения из партиции.

desired_partition_process_time

Тип: TDuration
Значение по умолчанию: 10m
Менять не рекомендуется. Регулирует, насколько крупные партиции будет нарезать контроллер.

desired_partition_rows_per_second

Тип: double
Значение по умолчанию: 1000.0
Менять не рекомендуется. Регулирует, насколько крупные партиции будет нарезать контроллер.

desired_partition_bytes_per_second

Тип: double
Значение по умолчанию: 1000000.0
Менять не рекомендуется. Регулирует, насколько крупные партиции будет нарезать контроллер.

read_timeout

Тип: TDuration
Значение по умолчанию: 5m
Менять не рекомендуется. Пересоздает table reader, если тот возвращает пустой ответ в течение периода.

Запись в статические таблицы в порядке прихода

Класс синка: 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

Параметр

Описание

output_directory

Тип: NYT::NYPath::TRichYPath
Обязательный параметр
Директория для выходных статических таблиц. В её атрибуте @progress хранится прогресс доставки: владелец (pipeline, computation, sink id), общая последовательность таблиц и frontier каждой партиции. Каждому синку нужна собственная директория: если владелец чужой, синк падает и просит удалить атрибут вручную.

table_period

Тип: TDuration
Значение по умолчанию: 5m
Шаг непрерывной сетки таблиц. Пустая таблица для шага T создаётся без пропуска, только когда системный watermark входного потока известен и строго больше T + table_period. Пустые таблицы в будущем не создаются.

table_ttl

Тип: TDuration
Обязательный параметр
Время жизни выходной таблицы: она удаляется через table_ttl после своего table_timestamp (через Cypress expiration_time). Отношение table_ttl / table_period ограничено 40000 из-за лимита Cypress на число детей директории.

table_name_format

Тип: std::string
Значение по умолчанию: %Y-%m-%dT%H:%M:%SZ
Формат UTC-таймстемпа в имени таблицы. Обязан кодировать таймстемп без потерь (валидируется по round-trip), иначе имена разных слотов совпали бы.

data_weight_column

Тип: std::optional<std::string>
Опциональная колонка с пользовательским весом сообщения для лимита max_data_weight (тип int64 или uint64, неотрицательные значения). Без неё, а также при null в ней, весом считается размер сообщения.

Динамическая спека:

NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnector::TArrivalOrderTableSink>

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

max_row_count

Тип: long
Значение по умолчанию: 10000
Максимальное число сообщений в таблице. По достижении лимита синк переходит к следующему шагу сетки, в том числе в будущее.

max_data_weight

Тип: long
Значение по умолчанию: 1073741824
Максимальный суммарный вес сообщений в таблице. По достижении лимита синк переходит к следующему шагу сетки.

Дополнительные параметры

transaction_timeout

Тип: TDuration
Значение по умолчанию: 5m
Таймаут мастер-транзакции, атомарно создающей таблицу и продвигающей прогресс доставки.

retry_backoff

Тип: TDuration
Значение по умолчанию: 1s
Задержка между повторными попытками транзакции. Попытки продолжаются до успеха или отмены job.

Импорт статической таблицы в стейт

Типичный сценарий: батч-процесс (YQL, …) периодически пересобирает справочник в виде статической YT-таблицы, а реалтайм-пайплайн должен джойнить свой поток против свежей версии этого справочника. Коннектор static_table поддерживает этот сценарий "из коробки", позволяя загрузить такую таблицу в стейт пайплайна с произвольной бизнес-логикой над каждой строкой перед записью.

Как это работает

Схема: один computation читает справочник через static_table-сорс и пишет его в стейт (loader); второй computation на входном потоке достаёт значение по ключу из того же стейта и эмитит обогащённое сообщение (enricher).

  1. static_table-сорс потоково отдаёт строки таблицы как обычные сообщения. Контроллер сам определяет, какая версия таблицы свежая, и подаёт её строки в пайплайн в виде отдельного стрима.
  2. Loader-computation (stateful) принимает каждую строку, при необходимости валидирует/нормализует её (trim, lower-case, обогащение из других источников, фильтрация и т. п.) и пишет результат в стейт по ключу.
  3. 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 и интеграционные тесты — в примерах.

Примеры

Полностью рабочие пайплайны с тестами, демонстрирующие данный подход:

Альтернативы

static_table расширение → стейт — не единственный способ организовать джойн со справочником; альтернативы и их компромиссы — в таблице ниже.

Подход

Когда выбирать

Цена

#1 static_table расширение → стейт (выше)

нужна бизнес-логика над строкой до записи в стейт; допустима полная перезаливка

стоимость перезаливки внутри пайплайна на каждой версии

#2 Конвертация в динамическую таблицу + symlink под external state

большой объём, lookup по запросу, нужна атомарная смена версии данных

сетевой lookup в динтаблицу + эксплуатация симлинка

#3 Встроенная БД и деплой через Resource

ноль сетевых обращений на джойне, объём ограничен памятью/диском воркера

сложная эксплуатация, формат и доставка БД на воркеры

В этом подходе роль стейта играет сама динамическая таблица: батч-процесс собирает новую версию справочника как сортированную динамическую таблицу …/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.

Примеры

В тестах примеров yt_sync сначала собирает reference.v1 и линкует current → reference.v1, пайплайн обогащает событие значением v1, затем собирается reference.v2, симлинк атомарно перенаправляется на новую версию, и для того же ключа пайплайн начинает отдавать значение v2.

Встроенная БД и деплой через Resource

Если справочник целиком помещается в память или на локальный SSD воркера, и доступ к нему нужен совсем без сетевых походов, имеет смысл собирать его как встроенную БД: батч-процесс выкатывает готовый файл/директорию, а пайплайн поднимает её как локальный лукап-движок прямо внутри процесса джоба.

Достоинства — нулевые сетевые обращения на джойне, минимальная задержка лукапа (только локальный диск/память), полная независимость от внешних KV-сервисов. Недостатки — объём ограничен ресурсами воркера; формат БД, её схема и совместимость версий ложатся на команду; доставку новой версии (Resource, переподготовка, чек целостности) нужно эксплуатировать самостоятельно. Подходит для относительно небольших справочников (единицы гигабайт), которые обновляются нечасто.

См. также

Предыдущая
Следующая