- Зачем: - необходимо расширить курс до silver-слоя с воспроизводимыми трансформациями. - Что: - создан файл plans/module-05-silver-layer.md с планом модуля. - обновлен plans/README.md с записью о новом модуле в таблицу содержания. - Проверка: - проверка согласованности ссылок и структуры документации.
27 KiB
Модуль 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-сессия останавливается.
План работ
- Создать
notebooks/05_silver_layer.ipynbпо ячеечной структуре, описанной ниже. - Обновить
plans/README.md— добавить строку Модуля 5 в таблицу оглавления. - Валидация: убедиться, что ноутбук выполняется сверху вниз на поднятом стенде после прохождения Модуля 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. Namespacelakehouse.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] Вопросы:
- Чем silver отличается от bronze?
- Зачем фиксировать правила перед кодом?
- Какие три уровня проверок качества и зачем?
- Почему LEFT JOIN, а не INNER JOIN с lookup?
- Что если фильтр удалит 50% строк?
- Можно ли пересоздать 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).