Разбираем, как из сырых метаданных dbt собирается аналитическая витрина: как определяется, какая версия записи «побеждает», что попадает в первый слой «как пришло», как из двух независимых слоёв — parse и compile — собирается общая модель узлов и почему отсутствующее значение — это NULL или пустой список, а не пропущенная строка.
Что определяет принадлежность записи к периоду
Каждая запись в витрине привязана к эпохеотдельному parquet-файлу с номером версии в имени; файлов много, потому что каждый прогон ingest добавляет новую эпоху, а старые остаются на диске. Файлы называются {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 — наборе идентификаторов живых узловузлов, которые присутствуют в текущем состоянии проекта; их список хранится в файле alive.parquet и читается по колонке unique_id. Отбор живых узлов встроен в тот же проход, что и дедупликация по эпохам, в функции 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способом выразить значение колонки при сборке строки: взять колонку из join-отношения, разобрать список из JSON или вычислить SQL-выражением. Вариант 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, потому что иначе скомпилированные поля не перенеслись бы вперёд. Остальные таблицы пишутся только при непустых строках.
Где смотреть в коде
- metadata_to_parquet.rs: trim_model_payload
- metadata_to_parquet.rs: dedup_epoch_groups
- parquet.rs: merge_compiled_nodes_into_nodes
- parquet.rs: new
- parquet.rs: write_dbt_table_merged
- metadata_to_parquet.rs: load_alive_ids
- metadata_to_parquet.rs: col_or_null
- metadata_to_parquet.rs: the_model_trim_keeps_what_the_row_builders_read