Назад к блогу

От исходных данных до аналитической витрины: цепочка решений, которую все недооценивают

От исходных данных до аналитической витрины: цепочка решений, которую все недооценивают

Статья разбирает редко обсуждаемый, но критичный участок пайплайна между сырыми метаданными dbt и готовой аналитической витриной. На примере конкретных решений — от выбора версии записи до логики дедупликации и отбора «живых» узлов — показано, как мелкие детали реализации определяют корректность итоговых данных. Полезно тем, кто строит или сопровождает подобные конвейеры и хочет понимать, где именно возникают неочевидные подводные камни.

Разбираем, как из сырых метаданных dbt собирается аналитическая витрина: как определяется, какая версия записи «побеждает», что попадает в первый слой «как пришло», как из двух независимых слоёв — parse и compile — собирается общая модель узлов и почему отсутствующее значение — это NULL или пустой список, а не пропущенная строка.

Что определяет принадлежность записи к периоду

Каждая запись в витрине привязана к эпохе. Файлы называются {SCHEMA_VERSION}_{N}.parquet, и номер N парсится из имени файла, после чего список эпох сортируется по этому номеру.

Принадлежность записи к периоду определяется именно номером эпохи, а не датами или статусами из payload. Причина в том, что порядок разрешения конфликтов (latest-wins) задаётся по номеру эпохи, а не по ingested_at: ingest упорядочивает записи по номеру эпохи, и две строки, записанные в одну секунду, иначе совпали бы по времени.

В SQL-выражении номер эпохи извлекается регулярным выражением из имени файла и приводится к BIGINT:

CAST(regexp_extract(filename, '_([0-9]+)\.parquet$', 1) AS BIGINT)

Приведение к BIGINT критично: при строковом сравнении v1_10 оказался бы ниже v1_9 (лексикографически '1' < '9'), тогда как числовое сравнение даёт 10 > 9. Комментарий к константе прямо указывает, что сортировка по числу, а не по строке, сохраняет v1_10 выше v1_9.

По этому номеру строится нумерация строк в оконных функциях. Режим Supersede::LatestBy оставляет по одной строке на ключ — самую новую по номеру эпохи. Режим Supersede::LatestGroup сохраняет все строки группы, относящиеся к самой новой эпохе, которая эту группу упомянула. Разница в том, что LatestBy выбирает одну строку на ключ, а LatestGroup — целую группу строк, объединённых ключом: если ключ описывает не одну запись, а набор связанных строк, LatestBy оставил бы из них только одну, а LatestGroup сохраняет весь набор из самой новой эпохи.

В Rust-пути ingest обрабатывает эпохи от новых к старым, отслеживая уже выданные unique_id, чтобы победила запись из более новой эпохи. При слиянии эпох победитель определяется по наибольшему ingested_at, а при равенстве — по большему номеру файла-эпохи; это подтверждает, что номер эпохи служит и тай-брейкером, и идентификатором версии файла.

Дедупликация и отбор живых узлов

Дедупликация идёт по группе, а не по строке. Группа — это весь набор строк одного узла (например, все колонки узла), и «побеждает» самая новая эпоха, которая эту группу упомянула. Эпохи читаются от новых к старым, и множество уже выданных идентификаторов (seen_ids) служит фильтром: если unique_id уже встречался в более новой эпохе, все его строки в старых эпохах отбрасываются.

Внутри одного батча идентификаторы собираются в отдельное множество batch_uids, и только после прохода по всем строкам батча они добавляются в seen_ids. Если помечать группу просмотренной сразу на первой строке, сохранилась бы только первая строка группы, а остальные потерялись бы.

Параллельно в том же проходе выполняется отбор живых узлов: строка отбрасывается, если её unique_id отсутствует в alive_ids — наборе идентификаторов живых узлов. Отбор живых узлов встроен в тот же проход, что и дедупликация по эпохам, в функции dedup_epoch_groups: она читает эпохи от новых к старым и для каждой строки проверяет, не мёртв ли узел и не видели ли его раньше, поэтому одна строка стоит одного поиска по хешу.

Ключ группировки задаётся параметром key_col. Для источников колонок compile/columns и catalog/columns это unique_id, где группа — набор колонок узла. Для compile/column_lineage это to_node_unique_id, где группа — входящие рёбра целевого узла.

Для parse/columns тот же принцип реализован отдельно: seen_ids помечается по эпохе, а prune_by_alive вызывается только при need_full.

Что сохраняется в первом слое «как пришло»

Payload модели не сохраняется целиком, потому что он много-килобайтный — в основном из-за raw_code, compiled_code и документации по колонкам. Вместо полного разбора через simd_json строится новый объект только из трёх полей.

config сохраняется как сырая JSON-строка, потому что её парсят только некоторые построители строк. Построитель, которому нужен config, сам разбирает эту строку: если payload содержит config как строку, она парсится в объект, а если config уже объект — берётся как есть.

__model_attr__ сохраняется целиком и разбирается, потому что это фиксированная небольшая структура без больших блобов кода и колонок. Это позволяет построителю, который начал читать новый атрибут, найти его здесь без изменения allowlist. Тест подтверждает: после trim config остаётся строкой {"materialized":"view"}, а в __model_attr__ сохраняются и contract, и primary_key, и time_spine, и даже other, который пока никто не читает.

Из __common_attr__ удаляется только raw_code. Остальные поля сохраняются, поэтому построителям строк остаются доступны patch_path, language, checksum и meta. Раньше эти поля отбрасывались целиком вместе со всем конвертом __common_attr__, из-за чего каждый model публиковал их как NULL.

Урезанная форма вызывается только для resource_type == "model". Для не-модельных узлов (macros, docs, sources и т.п.) делается полный simd_json-разбор, потому что у моделей в payload только config — нет raw_code, meta, patch_path и т.п., поэтому дорогой полный разбор для них избегается. Обрезка общая, чтобы delta-путь в dbt-index обрезал точно так же.

Первый слой: технические поля записи

В write_dbt_table, write_dbt_table_append и write_dbt_table_merged значение ingested_at проставляется одинаково: перед записью каждая строка приводится к объекту, и в него вставляется ключ "ingested_at" со строкой self.now, зафиксированной при создании IndexWriter.

Для типизированных строк (write_dbt_items) вставки нет: метод сразу вызывает write_table_items_typed, а структура строки уже должна содержать ingested_at последним полем. Причина в том, что serde_arrow обходит Serialize-реализацию напрямую без промежуточного serde_json::Value. В схеме таблицы поле ingested_at объявлено как не-null Timestamp(Microsecond, UTC) и идёт последним в списке полей, поэтому при типизированной записи значение должно быть уже в самой структуре, чтобы совпасть с последним полем схемы.

Момент фиксируется один раз в конструкторе IndexWriter: поле now заполняется вызовом chrono::Utc::now().to_rfc3339(), то есть строка времени создаётся при создании писателя и далее хранится в структуре. В cold_ingest_inner эта строка извлекается один раз и затем передаётся по ссылке во все функции записи таблиц одного прогона. Поэтому все строки, записанные через один экземпляр IndexWriter, получают одинаковую метку времени ingested_at, независимо от того, сколько времени заняла запись. Тот же &now передаётся и в write_parse_nodes, write_parse_columns, write_parse_project и остальные функции записи, что обеспечивает единую метку для всех таблиц прогона.

Второй слой: как parse и compile собираются в общую модель

parse/nodes и compile/nodes — это подкаталоги с эпохами в каталоге метаданных. parse/nodes читается как источник строк для dbt.nodes, а compile/nodes — как источник восьми полей, которые подмешиваются к этим строкам.

dbt.nodes строится как parse/nodes LEFT JOIN compile/nodes, потому что базовые строки узлов берутся из parse-эпох, а compile-эпохи лишь дополняют их. Compile-поля читаются заранее в CompileMap из колонок COMPILE_COLS (unique_id, compiled_code, compiled_code_hash, compiled_path, grain, grain_declared, grain_tested, classifiers, table_role) и затем подставляются в строки. Подстановка описывается типом EpochExpr. Вариант JoinCol берёт скалярное значение из присоединённой compile-записи, а JoinJsonList — список, закодированный в JSON-поле этой записи.

Восемь полей — compiled_code, compiled_code_hash, compiled_path, table_role, grain, grain_declared, grain_tested, classifiers — берутся из compile-записи через join-выражения: первые четыре как JoinCol, остальные четыре как JoinJsonList. Остальные колонки NodeRow берутся либо напрямую из epoch parquet, либо извлекаются из JSON-payload, либо жёстко заданы конверсией. Например, identifier вычисляется как COALESCE(t.identifier, json_extract_string(t.payload, '$.__source_attr__.identifier'), t.alias), а checksum — через JsonFirst(&["__common_attr__.checksum.checksum", "checksum.checksum"], "NULL").

Запись dbt.nodes идёт через write_dbt_items_merged в режиме CarryForwardMerge по ключу unique_id. В этом режиме переносятся те же compile-колонки: compiled_code, compiled_code_hash, compiled_path, grain, grain_declared, grain_tested, table_role.

Как compile-поля попадают в строки узлов

Компилируемые поля заранее собираются в карту CompileMap (тип HashMap<String, CompileFields>) функцией load_compile_nodes_map, которая читает все новые эпохи каталога compile/nodes и для каждой строки кладёт CompileFields под ключом unique_id. Затем merge_compiled_nodes_into_nodes проходит по уже существующим строкам dbt.nodes, берёт unique_id строки и, если он есть в карте, перезаписывает восемь полей через set_if_changed.

Если записи в карте нет, блок if let Some(c) = compile_map.get(&uid) не выполняется, и строка остаётся без изменений — это и есть семантика LEFT JOIN, сохраняющая parse-сторону: узел, отсутствующий в compile_map, остаётся ровно таким, каким был, поэтому частичный compile расширяет индекс, а не обрезает его.

set_if_changed вставляет значение только когда оно отличается от текущего, поэтому строки, у которых компилируемые колонки уже равны, не трогаются, и вызов, ничего не меняющий, не пишет файл. Это делает безопасным безусловный вызов после write_parse_nodes.

При построении строк узлов в NODES компилируемые колонки берутся через EpochExpr::JoinCol/JoinJsonList из compile/nodes через left join, а compiled_at, raw_code_hash и другие перечисленные поля жёстко заданы как EpochExpr::Null.

Обратный случай — compile-запись без parse-записи — в общую модель не попадает, потому что строки создаются из parse-батчей. write_parse_nodes читает только эпохи новее last_epoch_for(PARSE_NODES_SUBDIR) и строит node-строки из parse-батчей, а compile-поля лишь дополняют эти строки через compile_map. При отсутствии новых parse-эпох он не увидел бы новых эпох и ничего не записал, оставив compiled_code NULL на всю оставшуюся жизнь индекса. Именно поэтому вводится отдельный проход merge_compiled_nodes_into_nodes, который применяет compile-слой к уже лежащим на диске строкам, а не к конструируемым: это делает compile-слой применимым всякий раз, когда он появляется, а не только пока parse-эпохи ещё не израсходованы. В общей модели строки существуют только там, где их создал parse-проход, а compile-поля — это join-колонки, добавляемые к этим строкам; поэтому compile-запись без parse-записи не порождает новых строк и должна обрабатываться как UPDATE существующих, а не как построение.

Отсутствующее значение справочника

Отсутствующая запись представляется как NULL: при чтении значения из колонки проверяется флаг null, и если ячейка помечена как null, значение не извлекается.

Для строковых колонок это даёт NULL: в дельта-пути отсутствие строки обрабатывается через значение по умолчанию для необязательных полей и через пропуск строки для обязательных.

Для списочных колонок вместо NULL используется пустой список: в SQL-выражениях для не-nullable списков применяется COALESCE(t.fqn, []) и COALESCE(t.tags, []), потому что для null-колонки список читается как пустой. Схема это закрепляет: списочное поле объявлено non-nullable, тогда как строковое — nullable. Дополнительно для не-nullable списков строится all-empty (non-null) List<Utf8> array с нулевыми смещениями, то есть каждая строка получает пустой список, а не NULL.

Строка сохраняется потому, что соединение с compile-слоем делается как LEFT JOIN, а не INNER: при промахе поля остаются незаданными, а узел не отбрасывается. Когда файлов join-отношения нет на диске, его колонки деградируют: колоночные выражения становятся NULL, а списочные — пустым списком. На уровне конвертации отсутствие колонки в батче или несовпадение типа даёт null-массив нужной длины. Для не-nullable списков, которые должны быть пустыми, а не NULL, строится ListArray с n+1 нулевыми смещениями, так что каждый список имеет длину 0, а буфер значений пуст. В декларативной сборке колонок это же различие закреплено: для колонки без тела подставляется NULL, а для отсутствующего join-отношения списочное выражение заменяется на пустой список.

Пустая дельта всё равно пишет файл

Пустая дельта всё равно пишет файл nodes, потому что вызов записи nodes не обёрнут в проверку на непустоту, в отличие от остальных таблиц. Причина: даже при пустой дельте нужно перенести вперёд уже скомпилированные поля.

Для этого используется режим WriteMode::CarryForwardMerge с ключом unique_id и списком переносимых колонок: compiled_code, compiled_code_hash, compiled_path, grain, grain_declared, grain_tested, table_role. Остальные таблицы пишутся только если строки непусты — через макрос merge_if_nonempty!, который проверяет if !$rows.is_empty().

На дельта-ингесте (oldest-first, без inline-дедупликации) prune против alive_ids выполняется после записи: в write_parse_nodes при !iterate_newest_first (то есть когда need_full == false) сначала фильтруются node_rows и edge_rows по alive_ids, а затем для каждой специальной таблицы вызывается prune_by_alive. prune_by_alive оставляет только строки, чей unique_id присутствует в alive_ids, и ничего не делает, если alive_ids == None. Отличие от холодного ингеста в том, что при холодном ингесте iterate_newest_first = need_full = true, поэтому блок prune не выполняется, а вместо этого в дельта-ветке записи используется WriteMode::CarryForwardMerge с valid_ids — набором ключей, которые считаются живыми и потому должны пережить слияние: при слиянии сохраняются старые строки, чей ключ отсутствует в новом батче, но присутствует в valid_ids.

Что из этого следует на практике

Порядок версий определяется числом, извлечённым из имени файла, а не строкой и не временем записи. Любая сортировка, где номер эпохи остаётся строкой, даст неверный порядок на границе разрядности (v1_10 против v1_9), поэтому приведение к BIGINT — не косметика, а условие корректности.

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

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

Общая модель узлов строится из parse-строк, а compile-поля лишь дополняют их. Поэтому узел без compile-записи остаётся как есть, а compile-запись без parse-записи не порождает новой строки — её обрабатывает отдельный проход обновления уже лежащих на диске строк. Частичный compile расширяет индекс, а не обрезает его.

Отсутствующее значение — это NULL для строковых колонок и пустой список для списочных. Строка при этом не пропадает: LEFT JOIN, col_or_null и empty_list_array обеспечивают сохранение строки с NULL или пустым списком.

Пустая дельта всё равно пишет nodes, потому что иначе скомпилированные поля не перенеслись бы вперёд. Остальные таблицы пишутся только при непустых строках.

Где смотреть в коде

Источники

Похожее