Зачем: тикет #23 — каждое цитируемое число курса сверено с живым стендом; до этого в уроках стояли заглушки. Что: - уроки 00–04: 18 маркеров заполнены числами свежего импорта (280 437 событий, нули в таблицах ошибок, счётчики ODS/DDS); - лаба 07: таблица manifest после дня 4 (374 092 / 34 801 / 5 388), переходящие визиты по стыкам (35/22/26), числа после дня 5; - лаба 08: каноническая граница трёх дней, пример замера свежести (лаг 3:45 модельного времени до догона, ETL ~29 с) и вернувшегося пользователя; две живые поправки разбора времени: убран принудительный UTC в разборе STG и суффикс +00:00 в сравнении границы (ловились только на живом стенде). Проверка: детерминизм подтверждён двумя независимыми циклами сброс→импорт→инкремент (числа manifest и checksum_sha256 дней 4 и 5 совпали бит в бит); chain-check зелёный; сценарий лабы 08 прогнан вживую, включая стоп/продолжение и красный full_refresh=false из урока 4; grep «сверить-на-стенде» пуст; ссылки и якоря целы; make test (219+31) и make lint зелёные; /ai-text-lint по лабам — без существенных находок. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
26 KiB
Урок 2. STG → ODS: типизация и DQ-split
Формат: практика — будешь сам запускать команды и менять код, не только читать. Пререквизит: пройден урок 1 (слой STG — сырой JSON строкой уже лежит в
stg.*_raw, рядом метаданные доставки из Kafka). Эталонный путь:sql/ods/20_stg_to_ods.sqlи DDL целевых таблицsql/ddl/ods/20_ods.sql.Поток данных одной строкой:
stg.*_raw → ods.* (валидный ключ) + ods.*_errors (любая ошибка)О чём урок простыми словами: берём сырой JSON из STG, разбираем его на поля и приводим к типам, а заодно отделяем чистые записи от битых. И смотрим, что бывает, когда тип выбран неверно.
1. Зачем и где в проде
В прошлом уроке мы сложили сообщение в STG как есть — целым JSON-строкой. Никто его там не разбирал: задача STG была просто принять поток и ничего не уронить.
Теперь этот JSON пора разобрать. Каждое поле достаём из строки и приводим к нормальному типу:
event_id делаем UUID, время события — DateTime, координаты — числом. Зачем это нужно?
Пока значение лежит строкой, с ним почти ничего нельзя сделать: по строке не отфильтруешь
события за вчера, не сложишь координаты, не сравнишь числа. Как только поле стало настоящим
типом — с ним уже работают запросы. Слой, где данные впервые типизированы, и называется
ODS.
И ещё одно: именно здесь мы впервые начинаем отделять чистое от грязного. Поток никогда не бывает идеальным — где-то поле пустое, где-то вместо UUID мусор, где-то число записано как текст. Бросать такие записи нельзя (вдруг пригодятся для разбора), но и держать их вперемешку с чистыми — мешать себе же. Поэтому на входе в ODS поток раздваивается.
Правило, по которому всё раскладывается
Запомни его на весь урок — дальше всё держится на нём:
- у каждой таблицы есть ключ — поле, которое однозначно опознаёт запись. Для событий это
event_id, для контекста клика (устройство, гео) —click_id; - если ключ разобрался (получился валидным) — строка едет в основную таблицу
ods.*. Это «рабочие» данные, с которыми дальше живёт пайплайн; - а копия любой строки, где при разборе случилась хоть одна ошибка, едет в отдельную
таблицу ошибок
ods.*_errors. Туда складываем битое, чтобы потом разобрать, — и не теряем его, и не мешаем им чистым данным.
Вот это раздвоение по качеству и называется DQ-split (DQ — data quality, качество данных; split — разделение).
Осталась одна тонкость, к которой мы вернёмся в секции 3. Ошибка бывает не только в ключе.
Бывает, что ключ-то валидный, а испортилось какое-то другое поле. Тогда строка остаётся
в основной таблице (ключ на месте, она рабочая), но рядом, прямо в самой строке, ставится
пометка: вот это поле не разобралось. Пометки складываются в специальный столбец-список
parse_errors. И да — из-за этого одна запись может оказаться сразу в двух местах. Это не
ошибка, так задумано; почему — разберём ниже.
Почему батч, а не Materialized View
В уроке 1 остался открытый вопрос: поток в STG перекладывало Materialized View, почти в реальном времени, — почему дальше так не продолжить?
Разобрать STG → ODS через MV технически можно: оно бы типизировало каждое сообщение на лету,
по одному. Но мы сознательно идём другим путём — батчем. Батч значит вот что: всю таблицу
ODS мы пересобираем целиком, одной задачей Airflow. Сначала очищаем (TRUNCATE), потом
заново наполняем (INSERT из STG).
Зачем так, если MV быстрее? Ради двух вещей.
- Видно каждый прогон. Батч — это отдельная задача в Airflow: у неё есть запуск, статус, лог. Если что-то пошло не так, ты видишь, какой прогон сломался. MV же работает молча, фоном, и поймать момент сложнее.
- Пересчёт повторяем. Раз мы каждый раз чистим и наполняем заново, повторный запуск даёт ровно тот же результат. Захотел пересобрать слой — просто запусти задачу ещё раз.
На разборе типов и проверках качества это важнее, чем выиграть доли секунды на задержке. Это и есть «наблюдаемость и управляемость пересчёта» — одна из целей нашего стенда.
В проде иначе. Чистить и наполнять таблицу целиком каждый раз — это нормально для демо и маленького среза. На реальных объёмах так не делают: данные грузят инкрементально — добирают только новые, по «водяному знаку» (watermark — отметка, до какого момента уже всё загружено). Сама идея слоёв и DQ-split при этом не меняется.
2. Руки: смотрим базовый прогон
Подготовь стенд по
канонической инструкции курса.
После успешного world_init эталонный мир уже прошёл путь Kafka → STG → ODS → DDS → DM.
Ниже — форма блока «Статистика ODS», который печатает make transform. Это просто
счётчики строк по всем восьми таблицам слоя (четыре основных и четыре с ошибками):
Статистика ODS:
┌─table──────────────────────┬───rows─┐
│ ods.browser_event │ 280437 │
│ ods.location_event │ 280437 │
│ ods.device_by_click │ 26083 │
│ ods.geo_by_click │ 26083 │
│ ods.browser_event_errors │ 0 │
│ ods.location_event_errors │ 0 │
│ ods.device_by_click_errors │ 0 │
│ ods.geo_by_click_errors │ 0 │
└────────────────────────────┴────────┘
Прочитаем эту табличку — в ней три вещи, которые стоит заметить.
Все четыре *_errors — по нулям. Значит, стартовая история чистая: ни одна запись не дала
ошибки разбора, столбец parse_errors у всех пустой. Это нормально — данные стенда аккуратные.
Ошибки мы увидим в секции 4, когда сами их устроим.
browser и location идут в одном зерне события. Сколько событий пришло, столько строк
и ожидаем увидеть после типизации, если ключи валидны.
А device и geo обычно меньше, чем событий. Вот это уже интересно. Часть строк
куда-то делась? Нет. И это важно понять, иначе дальше будет казаться, что данные текут.
Дело в том, что эти две таблицы хранят не события, а контекст клика: с какого устройства
был клик и из какой точки на карте. Ключ у них — click_id. Разных кликов меньше, чем событий:
на один клик приходится несколько событий, и click_id у них повторяется. Движок таблицы
(про него — в секции 3) схлопывает повторы по ключу, оставляя по одной строке на клик.
Проверь это сам, а не верь на слово. Открой SQL-консоль http://localhost:9123/play
(пользователь default, пароль 123456) и посчитай, сколько в STG различных click_id:
-- Строк событий больше, чем различных click_id
SELECT count() AS stg_rows,
uniqExact(toUUIDOrNull(JSONExtractString(raw, 'click_id'))) AS distinct_clicks
FROM stg.geo_raw;
На эталонном мире запрос возвращает stg_rows = 280437 и
distinct_clicks = 26083. Второе число совпадает с числом строк в
ods.geo_by_click. Значит, это схлопнутые повторы по click_id, а не пропавшие
данные. Ничего не потерялось молча.
3. Загляни внутрь
Слой описан двумя файлами. Их полезно держать открытыми рядом — они про разное:
| Файл | Что задаёт |
|---|---|
sql/ddl/ods/20_ods.sql |
форму целевых таблиц: какие колонки, какие типы, какой движок |
sql/ods/20_stg_to_ods.sql |
наполнение: как из сырого JSON получить эти колонки |
Дальше — три места, ради которых урок и затевался. Пойдём по ним по порядку.
Типизация через *OrNull
Поле достаём из JSON и тут же приводим к нужному типу. Но не «жёстко», а через функции,
у которых на конце стоит OrNull:
toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id,
parseDateTime64BestEffortOrNull(JSONExtractString(raw, 'event_timestamp'), 6) AS event_ts,
toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude
Читается так: JSONExtractString(raw, 'event_id') достаёт поле из JSON как строку, а
toUUIDOrNull(...) пытается превратить эту строку в UUID.
Весь смысл — в суффиксе OrNull. Если значение не приводится к нужному типу (вместо UUID
пришёл мусор), функция не падает с ошибкой, а просто возвращает NULL. Это ровно то правило
стенда, что и в STG — «грязная запись не валит пайплайн», — только теперь на уровне типов.
Один кривой event_id станет NULL и будет помечен, а остальные строки спокойно доедут.
Кстати, про
AS: эти строки живут в блокеWITHв начале запроса.WITH— это просто способ заранее посчитать значение и дать ему имя, чтобы ниже по запросу ссылаться на него коротко, по имени, а не повторять всю формулу. Имя задаётся черезAS.
Сборка parse_errors
Теперь — как собирается тот самый список пометок. Какие именно поля не разобрались, видно вот здесь:
arrayFilter(x -> x != '', [
if(event_id IS NULL, 'bad_event_id', ''),
if(event_ts IS NULL, 'bad_event_timestamp', ''),
if(click_id IS NULL, 'bad_click_id', '')
]) AS parse_errors
Разберём изнутри. Сначала строится список меток: на каждое поле — своя строка. Если поле
вышло NULL (не разобралось) — кладём метку вроде 'bad_event_id', иначе — пустую строку
''. Потом arrayFilter выкидывает из списка все пустые строки. Что осталось — и есть список
«что сломалось в этой записи», прямо в самой строке данных. У чистой записи он пустой.
Сам split — и почему запись бывает в двух местах
Теперь главное. Одни и те же строки STG раскладываются по двум INSERT — в основную таблицу
и в таблицу ошибок. Отличаются они условием WHERE:
-- в основную таблицу: берём строки с валидным ключом
... WHERE event_id IS NOT NULL;
-- в таблицу ошибок: берём строки, где есть хоть одна ошибка разбора
... WHERE length(parse_errors) > 0
AND (event_id IS NULL OR event_ts IS NULL OR click_id IS NULL);
Обрати внимание: эти два условия пересекаются, и это сделано нарочно. Представь строку, у
которой event_id валидный, а вот event_timestamp пришёл битый. Что с ней происходит:
- в основную таблицу она попадёт — ключ (
event_id) на месте, строка рабочая. Рядом вparse_errorsбудет стоять меткаbad_event_timestamp; - и в таблицу ошибок она тоже попадёт — ошибка-то в ней есть.
Одна запись — в двух местах. Это и есть «двойной учёт», и у каждой таблицы тут своя роль. Основная отвечает на вопрос «что у нас есть для работы» (и честно помечает, где в строке изъян). Таблица ошибок отвечает на другой вопрос — «что пришло битым и требует разбора». В самом файле это записано комментарием в шапке, в блоке «DQ-split».
Заметь на будущее. Логика разбора в файле продублирована: каждое поле типизируется дважды — один раз в
INSERTосновной таблицы, другой раз вINSERTтаблицы ошибок (у каждого свойWITHс теми же формулами). Для учебного файла так нагляднее, но есть цена: если поменять разбор только в одном из двух мест, они разойдутся. В секции 4 мы как раз этим воспользуемся — и увидим, чем грозит такой рассинхрон.
Движок: почему строк контекста меньше
И последнее место — строчка про движок основных таблиц:
ENGINE = ReplacingMergeTree(src_ingest_ts)
ORDER BY (click_id)
ReplacingMergeTree — это таблица, которая схлопывает строки с одинаковым ключом (ключ берётся
из ORDER BY), оставляя самую свежую по src_ingest_ts — времени загрузки в ODS. Вот она,
причина разницы из секции 2: у device и geo много строк с одинаковым click_id, и движок
оставляет по одной на клик.
4. Управляемая правка: сломай тип — поймай тихую потерю
Урок про типы — так давай намеренно ошибёмся типом и посмотрим, что будет. Это самый поучительный момент урока.
Возьмём координату geo_latitude — широту. Это дробное число, например -7.60361. Достаём мы
её через toFloat64OrNull — «привести к дробному числу». Заменим тип на целочисленный —
toInt64OrNull, «привести к целому». Для строки "-7.60361" целого числа не получится
(там точка, дробная часть), и функция вернёт NULL. То есть широта просто исчезнет.
Из секции 3 помним: разбор продублирован, поэтому правок будет две — в обоих INSERT
блока GEO EVENTS. Открой sql/ods/20_stg_to_ods.sql, найди оба вхождения и в каждом замени
функцию:
-- было:
toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude
-- стало:
toInt64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude
Пересобираем слой:
make transform
И смотрим на ту же «Статистику ODS». Таблица ошибок гео, которая была пустой, теперь полная:
│ ods.geo_by_click │ ≈ 26083 │
│ ods.geo_by_click_errors │ 280437 │ ← было 0
В таблице ошибок число точное. В ods.geo_by_click после фонового схлопывания
останется 26083 строки, но сразу после прогона число может быть больше.
А в самой основной таблице широта пропала — но не молча, рядом стоит метка:
SELECT click_id, geo_latitude, geo_longitude, parse_errors
FROM ods.geo_by_click
LIMIT 4;
┌─click_id─────┬─geo_latitude─┬─geo_longitude─┬─parse_errors─────────┐
│ cee12466-... │ ᴺᵁᴸᴸ │ -8.07257 │ ['bad_geo_latitude'] │
│ a8e39850-... │ ᴺᵁᴸᴸ │ 37.92792 │ ['bad_geo_latitude'] │
└──────────────┴──────────────┴───────────────┴──────────────────────┘
Вот теперь видно всё разом — и DQ-split, и «двойной учёт» из секции 3 вживую. Строки с валидным
click_id остались в основной таблице с пометкой bad_geo_latitude. И те же записи попали в
geo_by_click_errors. Долгота на месте, а широты больше нет: один неверный тип — и целое поле
потеряно по всему слою. Заметили это parse_errors и таблица ошибок — для того DQ-split и нужен.
Бывает и хуже — тихо, совсем без метки. Здесь нас спас суффикс
OrNull: неверный тип далNULL, аNULLмы умеем замечать (на него и сработалparse_errors). По-настоящему опасен другой случай — когда неверный тип успешно возвращает неправильное значение. НиNULL, ни ошибки, ни метки — всё «зелёное», а данные испорчены. Ровно так в уроке 1 и нашёлся баг: времяkafka_tsприводили черезtoInt64(...)от значения типаDateTime64, это молча срезало миллисекунды, и время по всему стенду уехало в1970-01-21. Ничто на это не указывало — поймали только прогоном на стенде. Мораль урока: тип выбирают осознанно, даже когда функция «не падает».
Верни как было. Откати обе правки — верни toFloat64OrNull в оба места. Если запутался,
проще одной командой откатить весь файл к версии из репозитория:
git checkout -- sql/ods/20_stg_to_ods.sql
make transform
После этого geo_by_click_errors снова 0, широта на месте. А если стенд совсем
«поплыл», пройди
канонический сброс.
5. Проверь себя
| Действие | Где смотреть | Что ожидать |
|---|---|---|
| базовый прогон | блок «Статистика ODS» | основные таблицы не пустые, все *_errors = 0 |
почему device/geo меньше событий |
запрос uniqExact(click_id) по stg.geo_raw |
число различных click_id совпадает с ods.geo_by_click |
| правка из секции 4 | блок «Статистика ODS» | ods.geo_by_click_errors прыгнул с 0 на ненулевое число |
| правка из секции 4 | SELECT geo_latitude, parse_errors FROM ods.geo_by_click |
широта NULL, в parse_errors — bad_geo_latitude |
6. Что должно получиться
После урока у тебя на руках — видимый результат (одно на выбор):
- скрин блока «Статистика ODS», где после правки
ods.geo_by_click_errorsушёл с0на ненулевое число; - либо выборка из
ods.geo_by_clickс пустой широтой и меткойbad_geo_latitudeрядом.
И проверь себя на словах — примерно эти вопросы всплывут на еженедельном созвоне:
- чем функции с суффиксом
OrNullудобнее «жёсткого» приведения типа; - почему слой ODS мы наполняем батчем, а не Materialized View, как STG;
- почему одна и та же строка может оказаться и в основной таблице, и в
*_errors.
Если на последнем вопросе запнёшься — вернись к секции 3 и посмотри на условия WHERE у двух
INSERT. Ответ там.
Мост к уроку 3
Данные теперь типизированы и разложены по качеству. Но ods.device_by_click и
ods.geo_by_click — это всё ещё отдельные кусочки про один клик: устройство в одной
таблице, гео в другой. В уроке 3 (ODS → DDS) мы соберём из них цельную сущность — dds.click
(клик сразу с устройством и гео) — и таблицу событий dds.event. И там же наткнёмся на первый
вопрос целостности: а что делать с событием, у которого нет своего клика? Такие «сироты»
(orphan) — тема следующего урока.