Files
mini-lakehouse-lab/plans/module-05-silver-layer.md
T
ddadmin 82d067d6c6 docs(plan): добавлен план модуля 5 по silver-слою
- Зачем:
  - необходимо расширить курс до silver-слоя с воспроизводимыми трансформациями.
- Что:
  - создан файл plans/module-05-silver-layer.md с планом модуля.
  - обновлен plans/README.md с записью о новом модуле в таблицу содержания.
- Проверка:
  - проверка согласованности ссылок и структуры документации.
2026-03-07 20:03:18 +03:00

290 lines
27 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Модуль 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).