diff --git a/plans/README.md b/plans/README.md index 4efd315..a1396db 100644 --- a/plans/README.md +++ b/plans/README.md @@ -17,3 +17,4 @@ | 2. Ментальная модель Lakehouse: storage, catalog, compute | [module-02-lakehouse-mental-model.md](./module-02-lakehouse-mental-model.md) | Ready for validation | `notebooks/02_lakehouse_mental_model.ipynb` | 2026-03-07 | | 3. Raw-данные и первое чтение в Spark | [module-03-raw-ingest-and-first-read.md](./module-03-raw-ingest-and-first-read.md) | Ready for validation | `docker-compose.yml`, `.gitignore`, `START_HERE.md`, `docs/stack_reference.md`, `notebooks/03_raw_ingest_and_first_read.ipynb` | 2026-03-07 | | 4. Первая рабочая Iceberg-таблица и слой bronze | [module-04-bronze-with-iceberg.md](./module-04-bronze-with-iceberg.md) | Draft | `notebooks/04_bronze_with_iceberg.ipynb` | 2026-03-07 | +| 5. Слой silver и воспроизводимые трансформации | [module-05-silver-layer.md](./module-05-silver-layer.md) | Draft | `notebooks/05_silver_layer.ipynb` | 2026-03-07 | diff --git a/plans/module-05-silver-layer.md b/plans/module-05-silver-layer.md new file mode 100644 index 0000000..c6ae736 --- /dev/null +++ b/plans/module-05-silver-layer.md @@ -0,0 +1,289 @@ +# Модуль 5. Слой silver и воспроизводимые трансформации + +**Статус:** `Draft` +**Последнее обновление:** `2026-03-07` + +## Цель + +Научить студента переходу от bronze-таблицы «как есть» к осознанно очищенной и обогащённой silver-таблице. Студент формулирует правила трансформаций явно (не зарытыми в код), применяет их к bronze, обогащает таблицу вычисляемыми колонками и join-данными, а затем проверяет результат на трёх уровнях: строки, схема, агрегаты. + +## Результат для студента + +После прохождения модуля студент: + +- понимает, зачем нужен silver-слой и чем он отличается от bronze; +- умеет сформулировать правила трансформации явно и зафиксировать их как «реестр правил» перед кодом; +- умеет отфильтровать некачественные записи из bronze по осознанным критериям; +- умеет нормализовать типы колонок (DOUBLE -> INT для passenger_count, RatecodeID); +- умеет обработать NULL-значения явным образом (coalesce); +- умеет добавить вычисляемую колонку (trip_duration_minutes); +- умеет обогатить данные через JOIN с lookup-таблицей (названия зон pickup/dropoff); +- умеет выполнить проверки качества результата на трёх уровнях: строки, схема, агрегаты; +- может сопоставить bronze/silver с привычным ods/dds-мышлением. + +## Deliverables + +- `notebooks/05_silver_layer.ipynb` — основной ноутбук Модуля 5; +- обновление `plans/README.md` — строка Модуля 5 в оглавлении. + +Инфраструктурные изменения не требуются. + +## Зависимости + +### Предусловия + +- Стенд поднят и работает (Модуль 1). +- Студент прошёл Модуль 4: таблицы `lakehouse.bronze.nyc_taxi_yellow` и `lakehouse.bronze.taxi_zone_lookup` существуют. + +### Что используют последующие модули + +- Модуль 6 читает `lakehouse.silver.nyc_taxi_yellow` через Spark и Trino. +- Модуль 8 использует silver в финальной практике. + +## Дизайн-решения + +### Namespace: `lakehouse.silver` + +Семантическое имя (как `lakehouse.bronze` в Модуле 4). Подкрепляет мышление слоями, Модуль 6 будет читать `lakehouse.silver.*`. + +### Правила трансформаций: явный реестр в markdown перед кодом + +Ключевая методическая идея модуля — правила не должны быть зарыты в код. Перед блоком кода стоит markdown-таблица с реестром правил (аналог mapping specification в DWH-практике). Это делает трансформации прозрачными и подкрепляет PRD Risk #4. + +### Состав трансформаций: три группы + +**Группа 1 — Фильтрация:** +- F1: `fare_amount < 0` — отрицательные тарифы (аномалия из Модуля 3) +- F2: `total_amount > 1000` — экстремальные суммы (аномалия из Модуля 3) +- F3: `tpep_pickup_datetime` вне `[2024-01-01, 2025-01-01)` — записи за пределами учебного диапазона (аномалия из Модуля 3) + +**Группа 2 — Нормализация и обработка NULL:** +- N1: `passenger_count` — DOUBLE -> INT, NULL -> 0 +- N2: `RatecodeID` — DOUBLE -> INT, NULL -> 99 (unknown) +- N3: `congestion_surcharge` — NULL -> 0.0 +- N4: `Airport_fee` — NULL -> 0.0 + +**Группа 3 — Обогащение:** +- E1: `trip_duration_minutes` — вычисляемая колонка `(unix_timestamp(dropoff) - unix_timestamp(pickup)) / 60` +- E2: `pickup_borough`, `pickup_zone` — LEFT JOIN с taxi_zone_lookup по PULocationID +- E3: `dropoff_borough`, `dropoff_zone` — LEFT JOIN с taxi_zone_lookup по DOLocationID + +Ряд правил намеренно оставлен для самостоятельного задания: фильтрация `trip_distance = 0 AND total_amount != 0`, нормализация `payment_type` (BIGINT -> INT с именованными значениями), вычисляемая колонка `speed_mph`. + +### Диапазон дат для F3 + +Фиксированный `[2024-01-01, 2025-01-01)` — полный 2024 год, а не только каноническое подмножество `2024-01..2024-03` из `START_HERE.md`. Причина: фильтр ловит записи с явно ошибочным годом (2023, 2025 и т.д.), а не зависит от того, какие месяцы скачал студент. Канонический bundle — 3 месяца, расширенный — до 12, и в обоих случаях диапазон `[2024-01-01, 2025-01-01)` корректен. В Модуле 3 диапазон выводился из имён файлов через `infer_month_bounds`, но в silver мы работаем с bronze-таблицей (не с файлами), поэтому фиксированный диапазон проще и уместнее. + +### JOIN: два LEFT JOIN с алиасами + +LEFT JOIN (не INNER), чтобы не терять строки с LocationID без соответствия в lookup. Потеря строк из-за INNER JOIN была бы неявной трансформацией. + +### Создание таблицы: SQL CTAS с temp view + +Тот же паттерн, что в Модуле 4: трансформации через PySpark DataFrame API (`.filter()`, `.withColumn()`, `.join()`), затем `createOrReplaceTempView("silver_ready")`, затем `CREATE OR REPLACE TABLE lakehouse.silver.nyc_taxi_yellow USING iceberg AS SELECT col1, col2, ... FROM silver_ready` с явным списком колонок. + +### Проверки качества: три уровня + +- **Строковый (row-level):** assert нет fare_amount < 0, нет total_amount > 1000, нет NULL в passenger_count/RatecodeID +- **Схемный (schema-level):** наличие новых колонок, проверка типов (passenger_count = INT, RatecodeID = INT) +- **Агрегатный (aggregate-level):** таблица-сравнение bronze vs silver (row count, % отфильтрованных, средние fare/distance) + +### Taxi zone lookup: жёсткий prerequisite + +Lookup — самостоятельное задание Модуля 4. Модуль 5 работает только с bronze-таблицами и не должен лезть в raw-слой — это нарушило бы принцип разделения ответственности между слоями (PRD Risk #4). Стратегия: жёсткий assert при отсутствии `lakehouse.bronze.taxi_zone_lookup` с понятным сообщением и отсылкой к Модулю 4: «Таблица `lakehouse.bronze.taxi_zone_lookup` не найдена. Вернись в Модуль 4 и выполни самостоятельное задание (Секция 8).» + +### Cleanup: нет + +Silver нужен Модулям 6 и 8. Spark-сессия останавливается. + +## План работ + +1. Создать `notebooks/05_silver_layer.ipynb` по ячеечной структуре, описанной ниже. +2. Обновить `plans/README.md` — добавить строку Модуля 5 в таблицу оглавления. +3. Валидация: убедиться, что ноутбук выполняется сверху вниз на поднятом стенде после прохождения Модуля 4 (bronze-таблицы в каталоге). + +## Структура ноутбука + +### Секция 0: Введение +- **[md]** Заголовок, цели модуля, prerequisite (Модуль 4). Таблица-сравнение: + +| | Классический DWH (Greenplum) | Lakehouse | +|---|---|---| +| Источник | ods/stg-таблица | Bronze Iceberg-таблица | +| Результат | dds/fact-таблица с очищенными данными | Silver Iceberg-таблица | +| Правила трансформаций | Mapping specification / ETL-документация | Явный реестр правил в ноутбуке перед кодом | +| Проверки качества | dq-check скрипты, data quality framework | Asserts на трёх уровнях | +| Воспроизводимость | Повторный запуск ETL | `CREATE OR REPLACE TABLE ... AS SELECT` | + +- **[md]** Ключевая идея: silver = bronze + явные правила трансформации + проверки качества. + +### Секция 1: Spark-сессия и проверка bronze +- **[code]** SparkSession (паттерн Модулей 2-4: без `.master()`, `setLogLevel("ERROR")`). Константы, вспомогательные функции (`format_bytes`, `list_objects` — переиспользование паттерна из Модуля 4). +- **[code]** Assert: `spark.table(BRONZE_TABLE).count() > 0` с отсылкой к Модулю 4. Assert: `spark.table(BRONZE_LOOKUP_TABLE)` существует — при отсутствии жёсткая ошибка с отсылкой к самостоятельному заданию Модуля 4 (Секция 8). +- **[md]** Обе bronze-таблицы на месте. Теперь строим silver. + +### Секция 2: Профиль bronze — что нужно трансформировать +- **[md]** Прежде чем писать трансформации, понимаем текущее состояние данных. В Модуле 3 мы уже нашли аномалии. Теперь смотрим через bronze. +- **[code]** Краткий профиль: null-доли ключевых колонок, количество строк с fare_amount < 0, total_amount > 1000, pickup вне диапазона. +- **[md]** Таблица-профиль: какие проблемы нашли, сколько строк затронуто, что будем делать. + +### Секция 3: Реестр правил трансформации +- **[md]** Центральная идея модуля. Фиксируем все правила явно перед кодом. Аналог mapping specification в DWH. +- **[md]** Таблица правил (F1-F3, N1-N4, E1-E3) с колонками: #, Колонка/область, Правило, Тип, Обоснование. +- **[md]** Зачем это нужно: если через месяц спросят «почему в silver нет строк с отрицательным fare?», ответ в реестре, а не в git blame. + +### Секция 4: Применение правил — фильтрация +- **[md]** Начинаем с фильтрации. Порядок осознанный: сначала фильтруем, потом нормализуем и обогащаем (не тратим ресурсы на трансформацию строк, которые будут удалены). +- **[code]** Чтение bronze DataFrame. Применение F1, F2, F3. Подсчёт строк до/после, вывод количества и процента удалённых. +- **[md]** Обычно фильтры качества удаляют небольшую долю строк. Если фильтр убирает заметную часть данных — это сигнал перепроверить правило или источник. + +### Секция 5: Применение правил — нормализация +- **[md]** Приводим типы и обрабатываем NULL. Параллель с PostgreSQL: как `ALTER COLUMN TYPE integer`, но между слоями, а не in-place. +- **[code]** Применение N1-N4: `coalesce` + `cast` для passenger_count и RatecodeID, `coalesce` для surcharge и fee. Каждый `.withColumn()` с комментарием, какому правилу соответствует. + +### Секция 6: Применение правил — обогащение +- **[md]** Добавляем вычисляемые колонки и данные из lookup. LEFT JOIN — не теряем строки. Параллель с денормализацией в PostgreSQL. +- **[code]** E1: `trip_duration_minutes`. E2, E3: два LEFT JOIN с `taxi_zone_lookup` (алиасы `pu_zone`, `do_zone`). +- **[code]** Вывод нескольких строк с новыми колонками. + +### Секция 7: Создание silver-таблицы +- **[md]** Используем CTAS с `CREATE OR REPLACE`. Namespace `lakehouse.silver`. +- **[md]** Явное предупреждение о `CREATE OR REPLACE`: это полная перезапись таблицы с потерей предыдущего состояния и snapshot-истории. Для учебного silver это осознанный компромисс ради идемпотентности (повторный запуск ноутбука пересоздаёт таблицу с чистого листа). В рабочих сценариях это разрушительная операция — в Модуле 7 студент увидит, почему нужен другой подход. Параллель: `DROP TABLE + CREATE TABLE AS SELECT` в PostgreSQL. +- **[code]** `CREATE NAMESPACE IF NOT EXISTS lakehouse.silver`. Регистрация temp view. CTAS с явным списком демо-схемы: 24 колонки (19 оригинальных + 5 новых). Это каноническая схема, на которую опираются Модули 6 и 8. +- **[md]** Предупреждение о времени выполнения (2-4 мин на 3 мес данных). +- **[code]** Верификация: count, printSchema. + +### Секция 8: Проверки качества результата + +#### 8a: Строковый уровень +- **[code]** Assert: нет строк с fare_amount < 0, total_amount > 1000, NULL в passenger_count, NULL в RatecodeID. + +#### 8b: Схемный уровень +- **[code]** Проверка наличия новых колонок (trip_duration_minutes, pickup_borough, pickup_zone, dropoff_borough, dropoff_zone). Проверка типов: passenger_count = int, RatecodeID = int. + +#### 8c: Агрегатный уровень +- **[code]** Таблица-сравнение bronze vs silver: row count, % отфильтрованных, средние fare_amount, trip_distance. +- **[md]** Ориентир: если фильтры удалили больше 5% строк, это повод пересмотреть правила или проверить данные. + +### Секция 9: bronze vs silver — разделение ответственности +- **[md]** Мини-лекция (3-4 абзаца): raw -> bronze -> silver (-> gold, вне курса). Bronze = «как есть», silver = очищенные + обогащённые. Если правило неверно — bronze хранит оригинал, silver можно пересоздать. +- **[md]** Таблица-сравнение bronze vs silver по свойствам: данные, NULL, типы, доп. колонки, проверки, правила, «можно пересоздать из». + +### Секция 10: Самостоятельное задание +- **[md]** Инструкция: расширить silver-пайплайн тремя новыми правилами (по одному на каждый тип трансформации), пересоздать таблицу и проверить результат. + + **Фильтрация — F4:** удалить строки с `trip_distance = 0 AND total_amount != 0` (аномалия, найденная в Модуле 3, но не включённая в демонстрацию). + + **Нормализация — N5:** привести `payment_type` из BIGINT в INT и дать осмысленное значение для NULL (например, 0 = unknown). Подсказка: паттерн `coalesce` + `cast` тот же, что для `passenger_count`. + + **Обогащение — E4:** добавить вычисляемую колонку `speed_mph` — средняя скорость поездки в милях в час: `trip_distance / (trip_duration_minutes / 60)`. Подсказка: нужно обработать деление на ноль (если `trip_duration_minutes = 0`, результат должен быть NULL, а не ошибка). + + **Проверки:** после пересоздания silver написать: (1) row-level assert для F4, (2) schema-level проверку наличия `speed_mph` и типа `payment_type`, (3) сравнить количество строк до и после добавления F4. + +- **[code]** 6 пустых ячеек `# Ваш код: ...` (фильтрация, нормализация, обогащение, пересоздание таблицы, проверки, сравнение). + +### Секция 11: Checkpoint +- **[md]** Вопросы: + 1. Чем silver отличается от bronze? + 2. Зачем фиксировать правила перед кодом? + 3. Какие три уровня проверок качества и зачем? + 4. Почему LEFT JOIN, а не INNER JOIN с lookup? + 5. Что если фильтр удалит 50% строк? + 6. Можно ли пересоздать silver, если обнаружится новая аномалия? + +### Секция 12: Завершение +- **[md]** «Не удаляем silver-таблицу. В Модуле 6 будем читать через Spark и Trino, в Модуле 8 — финальная практика.» +- **[code]** `spark.stop()`. + +### Дизайн-решения по ноутбуку + +- **Helper-код:** весь код inline в ноутбуке. Студент должен видеть каждый шаг: фильтрацию, нормализацию, join, проверки. +- **Cleanup:** нет. Silver нужен Модулям 6 и 8. Spark-сессия останавливается. +- **Порядок трансформаций:** фильтрация -> нормализация -> обогащение. Логический порядок: сначала убираем мусор, потом чистим типы, потом добавляем контекст. +- **Spark API:** трансформации через PySpark DataFrame API (`.filter()`, `.withColumn()`, `.join()`), запись через SQL CTAS (консистентно с Модулем 4). + +## Схемы silver-таблицы + +### Демо-схема (24 колонки) — после секций 4-7 + +Это каноническая схема, на которую опираются downstream-модули (6, 8). Она создаётся демонстрационным кодом ноутбука и является минимальной гарантированной схемой `lakehouse.silver.nyc_taxi_yellow`. + +``` +VendorID BIGINT (без изменений) +tpep_pickup_datetime TIMESTAMP (без изменений) +tpep_dropoff_datetime TIMESTAMP (без изменений) +passenger_count INT (N1: DOUBLE->INT, NULL->0) +trip_distance DOUBLE (без изменений) +RatecodeID INT (N2: DOUBLE->INT, NULL->99) +store_and_fwd_flag STRING (без изменений) +PULocationID BIGINT (без изменений) +DOLocationID BIGINT (без изменений) +payment_type BIGINT (без изменений) +fare_amount DOUBLE (без изменений, F1 отфильтровал < 0) +extra DOUBLE (без изменений) +mta_tax DOUBLE (без изменений) +tip_amount DOUBLE (без изменений) +tolls_amount DOUBLE (без изменений) +improvement_surcharge DOUBLE (без изменений) +total_amount DOUBLE (без изменений, F2 отфильтровал > 1000) +congestion_surcharge DOUBLE (N3: NULL->0.0) +Airport_fee DOUBLE (N4: NULL->0.0) +trip_duration_minutes DOUBLE (E1: новая, вычисляемая) +pickup_borough STRING (E2: новая, из lookup JOIN) +pickup_zone STRING (E2: новая, из lookup JOIN) +dropoff_borough STRING (E3: новая, из lookup JOIN) +dropoff_zone STRING (E3: новая, из lookup JOIN) +``` + +### Расширенная схема (26 колонок) — после самостоятельного задания (секция 10) + +Студент пересоздаёт silver-таблицу с дополнительными правилами. Итоговая таблица содержит 24 демо-колонки + 2 новых, с изменением типа `payment_type`. Downstream-модули (6, 8) не зависят от расширенной схемы: `payment_type INT` обратно совместим с `BIGINT` для чтения, а дополнительные колонки (`speed_mph`) не используются downstream. + +``` +payment_type INT (N5: BIGINT->INT, NULL->0; было BIGINT) +speed_mph DOUBLE (E4: новая, trip_distance / (trip_duration_minutes / 60)) +``` + +Также добавляется фильтр F4 (`trip_distance = 0 AND total_amount != 0`), который уменьшает количество строк. + +## Checkpoint + +Студент должен уметь: + +- создать silver-таблицу из bronze с явными правилами трансформации; +- показать реестр правил и объяснить обоснование каждого; +- показать результаты проверок качества на трёх уровнях; +- объяснить, чем silver отличается от bronze; +- объяснить, почему silver можно пересоздать из bronze. + +## Acceptance Criteria + +- `notebooks/05_silver_layer.ipynb` можно выполнить сверху вниз в поднятом Jupyter-окружении после прохождения Модуля 4 (bronze-таблицы в каталоге). +- Ноутбук следует методике курса: `объяснение -> демонстрация -> самостоятельное повторение -> checkpoint`. +- Ноутбук использует паттерны из Модулей 1-4: SparkSession без `.master()`, `setLogLevel("ERROR")`. +- Ноутбук идемпотентен: повторное выполнение (`CREATE OR REPLACE`) не дублирует данные и не роняет ячейки. +- CTAS использует явный список колонок, а не `SELECT *`. +- Правила трансформации зафиксированы в markdown-таблице перед кодом, а не только в коде. +- Проверки качества выполняются на трёх уровнях: строки, схема, агрегаты. +- Silver-таблица `lakehouse.silver.nyc_taxi_yellow` остаётся после выполнения ноутбука для использования в Модулях 6 и 8. +- При отсутствии bronze-таблиц — понятный assert с отсылкой к Модулю 4. +- Текст ноутбука на русском языке с параллелями к PostgreSQL/Greenplum. +- `plans/README.md` содержит строку Модуля 5 в оглавлении. + +## Риски + +- **Taxi zone lookup отсутствует.** Самостоятельное задание Модуля 4 может быть не выполнено. Жёсткий assert с отсылкой к Модулю 4. Модуль 5 не создаёт bronze-таблицы самостоятельно — это нарушило бы разделение ответственности между слоями. +- **Время CTAS.** JOIN + трансформации — на 30-50% медленнее bronze CTAS. На 3 мес (~7 млн строк): 2-4 мин. Extended: до 5-10 мин. Предупреждение в markdown. +- **Column case sensitivity.** Новые колонки (lowercase) рядом с оригинальными (mixed case: VendorID, PULocationID). Нормально для Iceberg, но добавить пояснение. +- **trip_duration_minutes отрицательные.** Если dropoff < pickup — длительность отрицательна. Упомянуть как наблюдение, не фильтровать в демо (потенциально — для самостоятельного задания или будущего правила). +- **Идемпотентность.** `CREATE OR REPLACE TABLE ... AS SELECT` (CTAS) для повторных запусков — тот же паттерн, что в Модуле 4. Пересоздаёт таблицу целиком. + +## Out of Scope + +- Gold-слой и бизнес-витрины. +- Чтение silver через Trino — отложено до Модуля 6. +- Сложные трансформации (MERGE, SCD, window functions). +- dq-фреймворки (Great Expectations, dbt tests) — принцип показан через assert. +- Time travel и schema evolution — отложено до Модуля 7. +- Compaction и vacuum — отложено до Модуля 8. +- Партиционирование silver-таблицы (PRD: вне v1).