diff --git a/docs/course/lessons/02_stg_to_ods.md b/docs/course/lessons/02_stg_to_ods.md new file mode 100644 index 0000000..5c0f39a --- /dev/null +++ b/docs/course/lessons/02_stg_to_ods.md @@ -0,0 +1,362 @@ +# Урок 2. STG → ODS: типизация и DQ-split + +> Статус: черновик. Режим: **руки**. +> Пререквизит: пройден урок 1 (слой STG — сырой JSON строкой уже лежит в `stg.*_raw`, +> рядом метаданные доставки из Kafka). +> Эталонный путь: [`sql/ods/20_stg_to_ods.sql`](../../../sql/ods/20_stg_to_ods.sql) +> и DDL целевых таблиц [`sql/ddl/ods/20_ods.sql`](../../../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. Руки: смотрим базовый прогон + +Поднимаем стенд, создаём схему, заливаем **малый срез** (50 строк на топик) и запускаем +трансформацию: + +```bash +make up # поднять инфраструктуру +make ddl # создать базы и таблицы (в т.ч. слой ODS) +LIMIT=50 make data # залить по 50 строк каждого файла в Kafka → STG +make transform # батч STG → ODS → DDS → DM (нас интересует первый шаг) +``` + +`make transform` прогоняет всю цепочку слоёв сразу, но прямо в консоли печатает то, что нам +нужно сейчас, — блок **«Статистика ODS»**. Это просто счётчики строк по всем восьми таблицам +слоя (четыре основных и четыре с ошибками): + +``` +Статистика ODS: + ┌─table──────────────────────┬─rows─┐ + │ ods.browser_event │ 50 │ + │ ods.browser_event_errors │ 0 │ + │ ods.location_event │ 50 │ + │ ods.location_event_errors │ 0 │ + │ ods.device_by_click │ 26 │ + │ ods.device_by_click_errors │ 0 │ + │ ods.geo_by_click │ 26 │ + │ ods.geo_by_click_errors │ 0 │ + └────────────────────────────┴──────┘ +``` + +Прочитаем эту табличку по строчкам — в ней три вещи, которые стоит заметить. + +**Все четыре `*_errors` — по нулям.** Значит, наш срез чистый: ни одна запись не дала ошибки +разбора, столбец `parse_errors` у всех пустой. Это нормально — данные в демо аккуратные. +Ошибки мы увидим в секции 4, когда сами их устроим. + +**`browser` и `location` дали 50 из 50.** Сколько событий пришло — столько и легло, один к +одному. + +**А `device` и `geo` — только 26 из 50.** Вот это уже интересно. Половина куда-то делась? Нет. +И это важно понять, иначе дальше будет казаться, что данные текут. + +Дело в том, что эти две таблицы хранят не события, а **контекст клика**: с какого устройства +был клик и из какой точки на карте. Ключ у них — `click_id`. А в срезе на 50 событий разных +кликов всего 26: на один клик приходится несколько событий, и `click_id` у них повторяется. +Движок таблицы (про него — в секции 3) схлопывает повторы по ключу, оставляя по одной строке +на клик. Отсюда и 26. + +Проверь это сам, а не верь на слово. Открой SQL-консоль `http://localhost:9123/play` +(пользователь `default`, пароль `123456`) и посчитай, сколько в срезе *различных* `click_id`: + +```sql +-- Всего строк в STG — 50, но различных click_id среди них — ровно 26 +SELECT count() AS stg_rows, + uniqExact(toUUIDOrNull(JSONExtractString(raw, 'click_id'))) AS distinct_clicks +FROM stg.geo_raw; +``` + +Получишь `stg_rows = 50`, `distinct_clicks = 26` — ровно столько, сколько строк в +`ods.geo_by_click`. Значит, 26 — это схлопнутые повторы, а не пропавшие данные. Ничего не +потерялось молча. + +--- + +## 3. Загляни внутрь + +Слой описан **двумя файлами**. Их полезно держать открытыми рядом — они про разное: + +| Файл | Что задаёт | +|------|------------| +| `sql/ddl/ods/20_ods.sql` | **форму** целевых таблиц: какие колонки, какие типы, какой движок | +| `sql/ods/20_stg_to_ods.sql` | **наполнение**: как из сырого JSON получить эти колонки | + +Дальше — три места, ради которых урок и затевался. Пойдём по ним по порядку. + +### Типизация через `*OrNull` + +Поле достаём из JSON и тут же приводим к нужному типу. Но не «жёстко», а через функции, +у которых на конце стоит `OrNull`: + +```sql +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` и будет помечен, а остальные 49 строк спокойно доедут. + +> Кстати, про `AS`: эти строки живут в блоке `WITH` в начале запроса. `WITH` — это просто +> способ заранее посчитать значение и дать ему имя, чтобы ниже по запросу ссылаться на него +> коротко, по имени, а не повторять всю формулу. Имя задаётся через `AS`. + +### Сборка `parse_errors` + +Теперь — как собирается тот самый список пометок. Какие именно поля не разобрались, видно вот +здесь: + +```sql +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`: + +```sql +-- в основную таблицу: берём строки с валидным ключом +... 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 мы как раз этим +> воспользуемся — и увидим, чем грозит такой рассинхрон. + +### Движок: откуда взялись 26 строк + +И последнее место — строчка про движок основных таблиц: + +```sql +ENGINE = ReplacingMergeTree(src_ingest_ts) +ORDER BY (click_id) +``` + +`ReplacingMergeTree` — это таблица, которая схлопывает строки с одинаковым ключом (ключ берётся +из `ORDER BY`), оставляя самую свежую по `src_ingest_ts` — времени загрузки в ODS. Вот она, +причина «26 из 50» из секции 2: у `device` и `geo` много строк с одинаковым `click_id`, и +движок оставляет по одной на клик. + +--- + +## 4. Управляемая правка: сломай тип — поймай тихую потерю + +Урок про типы — так давай **намеренно ошибёмся типом** и посмотрим, что будет. Это самый +поучительный момент урока. + +Возьмём координату `geo_latitude` — широту. Это дробное число, например `50.82709`. Достаём мы +её через `toFloat64OrNull` — «привести к дробному числу». Заменим тип на целочисленный — +`toInt64OrNull`, «привести к целому». Для строки `"50.82709"` целого числа не получится +(там точка, дробная часть), и функция вернёт `NULL`. То есть широта просто исчезнет. + +Из секции 3 помним: разбор продублирован, поэтому правок будет **две** — в обоих `INSERT` +блока `GEO EVENTS`. Открой `sql/ods/20_stg_to_ods.sql`, найди оба вхождения и в каждом замени +функцию: + +```sql +-- было: +toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude +-- стало: +toInt64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude +``` + +Пересобираем слой: + +```bash +make transform +``` + +И смотрим на ту же «Статистику ODS». Таблица ошибок гео, которая была пустой, теперь полная: + +``` + │ ods.geo_by_click │ 26 │ + │ ods.geo_by_click_errors │ 50 │ ← было 0 +``` + +А в самой основной таблице широта пропала — но не молча, рядом стоит метка: + +```sql +SELECT click_id, geo_latitude, geo_longitude, parse_errors +FROM ods.geo_by_click +LIMIT 4; +``` + +``` +┌─click_id─────┬─geo_latitude─┬─geo_longitude─┬─parse_errors─────────┐ +│ 58cdfc1e-... │ ᴺᵁᴸᴸ │ -0.2 │ ['bad_geo_latitude'] │ +│ 9ffd819b-... │ ᴺᵁᴸᴸ │ 85.37752 │ ['bad_geo_latitude'] │ +└──────────────┴──────────────┴───────────────┴──────────────────────┘ +``` + +Вот теперь видно всё разом — и DQ-split, и «двойной учёт» из секции 3 вживую. 26 строк +остались в основной таблице (ключ `click_id` цел) с пометкой `bad_geo_latitude`. И те же +записи попали в число 50 строк `geo_by_click_errors`. Долгота на месте, а широты больше нет: +один неверный тип — и целое поле потеряно по всему слою. Заметили это `parse_errors` и таблица +ошибок — для того DQ-split и нужен. + +> **Бывает и хуже — тихо, совсем без метки.** Здесь нас спас суффикс `OrNull`: неверный тип +> дал `NULL`, а `NULL` мы умеем замечать (на него и сработал `parse_errors`). По-настоящему +> опасен другой случай — когда неверный тип **успешно** возвращает *неправильное* значение. +> Ни `NULL`, ни ошибки, ни метки — всё «зелёное», а данные испорчены. Ровно так в уроке 1 и +> нашёлся баг: время `kafka_ts` приводили через `toInt64(...)` от значения типа `DateTime64`, +> это молча срезало миллисекунды, и время по всему стенду уехало в `1970-01-21`. Ничто на это +> не указывало — поймали только прогоном на стенде. Мораль урока: тип выбирают осознанно, даже +> когда функция «не падает». + +**Верни как было.** Откати обе правки — верни `toFloat64OrNull` в оба места. Если запутался, +проще одной командой откатить весь файл к версии из репозитория: + +```bash +git checkout -- sql/ods/20_stg_to_ods.sql +make transform +``` + +После этого `geo_by_click_errors` снова `0`, широта на месте. А если стенд совсем «поплыл» — +всегда есть полный сброс: `make clean && make up && make ddl && LIMIT=50 make data && make transform`. + +--- + +## 5. Проверь себя + +| Действие | Где смотреть | Что ожидать | +|----------|--------------|-------------| +| `make transform` (базовый прогон) | блок «Статистика ODS» | `browser`/`location` = 50, `device`/`geo` = 26, все `*_errors` = 0 | +| почему 26, а не 50 | запрос `uniqExact(click_id)` по `stg.geo_raw` | 26 различных `click_id` — это схлопывание повторов, а не потеря | +| правка из секции 4 | блок «Статистика ODS» | `ods.geo_by_click_errors` прыгнул `0 → 50` | +| правка из секции 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` на `50`; +- либо выборка из `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) — тема следующего урока. diff --git a/sql/ddl/ods/20_ods.sql b/sql/ddl/ods/20_ods.sql index f9df01c..1493367 100644 --- a/sql/ddl/ods/20_ods.sql +++ b/sql/ddl/ods/20_ods.sql @@ -9,21 +9,9 @@ -- -- Важно: -- - Наполнение ODS выполняется batch-процессом из sql/ods/20_stg_to_ods.sql --- - Materialized View STG → ODS в target-state не используются +-- (а не Materialized View, как в STG): батч ради наблюдаемости пересчёта -- ============================================================================ --- ---------------------------------------------------------------------------- --- Удаляем legacy MV STG → ODS (если ранее были созданы) --- ---------------------------------------------------------------------------- -DROP TABLE IF EXISTS stg.mv_browser_raw_to_ods; -DROP TABLE IF EXISTS stg.mv_browser_raw_to_ods_errors; -DROP TABLE IF EXISTS stg.mv_location_raw_to_ods; -DROP TABLE IF EXISTS stg.mv_location_raw_to_ods_errors; -DROP TABLE IF EXISTS stg.mv_device_raw_to_ods; -DROP TABLE IF EXISTS stg.mv_device_raw_to_ods_errors; -DROP TABLE IF EXISTS stg.mv_geo_raw_to_ods; -DROP TABLE IF EXISTS stg.mv_geo_raw_to_ods_errors; - -- ============================================================================ -- BROWSER EVENTS -- ============================================================================ @@ -46,6 +34,8 @@ CREATE TABLE IF NOT EXISTS ods.browser_event parse_errors Array(LowCardinality(String)) -- Ошибки парсинга (если есть) ) ENGINE = ReplacingMergeTree(src_ingest_ts) +-- Партиционируем по бизнес-дате события: у browser_event есть собственное время (event_ts). +-- У click-контекста ниже (location/device/geo) такого времени нет — там партиция по дате загрузки. PARTITION BY toYYYYMM(event_date) ORDER BY (event_id) SETTINGS allow_nullable_key = 1; @@ -87,6 +77,7 @@ CREATE TABLE IF NOT EXISTS ods.location_event parse_errors Array(LowCardinality(String)) ) ENGINE = ReplacingMergeTree(src_ingest_ts) +-- По дате загрузки: у location нет собственного времени события (только время приёма в STG) PARTITION BY toYYYYMM(toDate(src_ingest_ts)) ORDER BY (event_id) SETTINGS allow_nullable_key = 1; diff --git a/sql/ods/20_stg_to_ods.sql b/sql/ods/20_stg_to_ods.sql index 50b4694..882c174 100644 --- a/sql/ods/20_stg_to_ods.sql +++ b/sql/ods/20_stg_to_ods.sql @@ -1,11 +1,20 @@ -- ============================================================================ -- Batch-трансформация: STG → ODS -- ============================================================================ +-- Поток данных: +-- stg.*_raw → ods.* (валидный ключ) + ods.*_errors (любая ошибка) +-- -- Назначение: -- - Перенос типизации STG → ODS из Materialized View в управляемый batch -- - Полная пересборка ODS для прозрачного мониторинга в Airflow -- - Сохранение DQ-логики: parse_errors + отдельные *_errors таблицы -- +-- DQ-split (почему одна строка может попасть в оба места): +-- - Основная таблица ods.* — строки с валидным КЛЮЧОМ (event_id / click_id). +-- В ней допустимы parse_errors по НЕключевым полям — строка остаётся, но помечена. +-- - Таблица ods.*_errors — копия строк с ЛЮБОЙ ошибкой парсинга (для разбора). +-- Поэтому строка с валидным ключом, но битым неключевым полем, попадёт И туда, И туда. +-- -- Когда запускать: -- - В DAG etl_pipeline перед ODS → DDS -- - После загрузки очередного среза данных в Kafka/STG