From bcf7975227e0bcd7c4bf417c4b33c6d875b320c5 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Mar 2026 23:32:43 +0300 Subject: [PATCH] =?UTF-8?q?feat(notebooks):=20=D1=80=D0=B5=D0=B0=D0=BB?= =?UTF-8?q?=D0=B8=D0=B7=D0=BE=D0=B2=D0=B0=D0=BD=20=D0=BC=D0=BE=D0=B4=D1=83?= =?UTF-8?q?=D0=BB=D1=8C=205=20=D0=BF=D0=BE=20=D1=81=D0=BB=D0=BE=D1=8E=20si?= =?UTF-8?q?lver=20=D0=B8=20=D1=82=D1=80=D0=B0=D0=BD=D1=81=D1=84=D0=BE?= =?UTF-8?q?=D1=80=D0=BC=D0=B0=D1=86=D0=B8=D1=8F=D0=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - обучение переходу от bronze-таблицы к очищенному и обогащенному silver-слою. - Что: - создан notebooks/05_silver_layer.ipynb с пайплайном фильтрации, нормализации и JOIN. - в ноутбук добавлены проверки качества на уровнях строк, схемы и агрегатов. - статус модуля 5 в планах обновлен до Ready for validation. - Проверка: - ручная проверка структуры ноутбука и корректности Spark-трансформаций. --- notebooks/05_silver_layer.ipynb | 627 ++++++++++++++++++++++++++++++++ plans/README.md | 2 +- plans/module-05-silver-layer.md | 2 +- 3 files changed, 629 insertions(+), 2 deletions(-) create mode 100644 notebooks/05_silver_layer.ipynb diff --git a/notebooks/05_silver_layer.ipynb b/notebooks/05_silver_layer.ipynb new file mode 100644 index 0000000..03f6d60 --- /dev/null +++ b/notebooks/05_silver_layer.ipynb @@ -0,0 +1,627 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "id": "section_0", + "metadata": {}, + "source": [ + "# Модуль 5. Слой silver и воспроизводимые трансформации\n", + "\n", + "В этом модуле мы переходим от bronze-таблицы «как есть» к осознанно очищенной и обогащённой silver-таблице.\n", + "\n", + "**Цели модуля:**\n", + "- понять, зачем нужен silver-слой и чем он отличается от bronze;\n", + "- научиться формулировать правила трансформации явно перед написанием кода;\n", + "- отфильтровать некачественные записи, нормализовать типы и обработать NULL;\n", + "- обогатить данные через JOIN с lookup-таблицей;\n", + "- выполнить проверки качества результата на трёх уровнях: строки, схема, агрегаты.\n", + "\n", + "**Prerequisite:**\n", + "- пройден Модуль 4: таблицы `lakehouse.bronze.nyc_taxi_yellow` и `lakehouse.bronze.taxi_zone_lookup` (самостоятельное задание) существуют.\n", + "\n", + "### Связка с привычным DWH (PostgreSQL / Greenplum)\n", + "\n", + "| | Классический DWH (Greenplum) | Lakehouse |\n", + "|---|---|---|\n", + "| Источник | ods/stg-таблица | Bronze Iceberg-таблица |\n", + "| Результат | dds/fact-таблица с очищенными данными | Silver Iceberg-таблица |\n", + "| Правила трансформаций | Mapping specification / ETL-документация | Явный реестр правил в ноутбуке перед кодом |\n", + "| Проверки качества | dq-check скрипты, data quality framework | Asserts на трёх уровнях |\n", + "| Воспроизводимость | Повторный запуск ETL | `CREATE OR REPLACE TABLE ... AS SELECT` |\n", + "\n", + "Ключевая идея: **silver = bronze + явные правила трансформации + проверки качества**." + ] + }, + { + "cell_type": "markdown", + "id": "section_1_title", + "metadata": {}, + "source": [ + "## 1. Spark-сессия и проверка bronze\n", + "\n", + "Инициализируем Spark и убеждаемся, что обе bronze-таблицы на месте. Теперь мы работаем только с ними и не обращаемся к raw-файлам напрямую." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_1_code", + "metadata": {}, + "outputs": [], + "source": [ + "import pandas as pd\n", + "from pyspark.sql import SparkSession, functions as F\n", + "\n", + "CATALOG_NAME = \"lakehouse\"\n", + "BRONZE_NAMESPACE = f\"{CATALOG_NAME}.bronze\"\n", + "BRONZE_TABLE = f\"{BRONZE_NAMESPACE}.nyc_taxi_yellow\"\n", + "BRONZE_LOOKUP_TABLE = f\"{BRONZE_NAMESPACE}.taxi_zone_lookup\"\n", + "\n", + "SILVER_NAMESPACE = f\"{CATALOG_NAME}.silver\"\n", + "SILVER_TABLE = f\"{SILVER_NAMESPACE}.nyc_taxi_yellow\"\n", + "\n", + "spark = SparkSession.builder \\\n", + " .appName(\"module-05-silver-layer\") \\\n", + " .getOrCreate()\n", + "\n", + "spark.sparkContext.setLogLevel(\"ERROR\")\n", + "\n", + "print(\"Spark готов\")\n", + "print(f\"Bronze table: {BRONZE_TABLE}\")\n", + "print(f\"Silver table: {SILVER_TABLE}\")" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_1_asserts", + "metadata": {}, + "outputs": [], + "source": [ + "assert spark.catalog.tableExists(BRONZE_TABLE), (\n", + " f\"Таблица {BRONZE_TABLE} не найдена. Вернись в Модуль 4 и выполни загрузку в bronze.\"\n", + ")\n", + "\n", + "assert spark.catalog.tableExists(BRONZE_LOOKUP_TABLE), (\n", + " f\"Таблица {BRONZE_LOOKUP_TABLE} не найдена. \"\n", + " \"Вернись в Модуль 4 и выполни самостоятельное задание (Секция 8).\"\n", + ")\n", + "\n", + "bronze_count = spark.table(BRONZE_TABLE).count()\n", + "assert bronze_count > 0, f\"Таблица {BRONZE_TABLE} пуста.\"\n", + "\n", + "print(f\"Обе bronze-таблицы на месте. В основной таблице {bronze_count:,} строк. Теперь строим silver.\")" + ] + }, + { + "cell_type": "markdown", + "id": "section_2_title", + "metadata": {}, + "source": [ + "## 2. Профиль bronze — что нужно трансформировать\n", + "\n", + "Прежде чем писать трансформации, понимаем текущее состояние данных. В Модуле 3 мы уже нашли аномалии. Теперь смотрим через bronze." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_2_code", + "metadata": {}, + "outputs": [], + "source": [ + "bronze_df = spark.table(BRONZE_TABLE)\n", + "\n", + "profile_df = bronze_df.select(\n", + " F.count(\"*\").alias(\"total_rows\"),\n", + " (F.sum(F.when(F.col(\"fare_amount\") < 0, 1).otherwise(0)) / F.count(\"*\") * 100).alias(\"negative_fare_pct\"),\n", + " (F.sum(F.when(F.col(\"total_amount\") > 1000, 1).otherwise(0)) / F.count(\"*\") * 100).alias(\"extreme_total_pct\"),\n", + " (F.sum(F.when((F.col(\"tpep_pickup_datetime\") < \"2024-01-01\") | (F.col(\"tpep_pickup_datetime\") >= \"2025-01-01\"), 1).otherwise(0)) / F.count(\"*\") * 100).alias(\"out_of_range_date_pct\"),\n", + " (F.sum(F.when(F.col(\"passenger_count\").isNull(), 1).otherwise(0)) / F.count(\"*\") * 100).alias(\"null_passengers_pct\"),\n", + " (F.sum(F.when(F.col(\"RatecodeID\").isNull(), 1).otherwise(0)) / F.count(\"*\") * 100).alias(\"null_ratecode_pct\"),\n", + " (F.sum(F.when(F.col(\"congestion_surcharge\").isNull(), 1).otherwise(0)) / F.count(\"*\") * 100).alias(\"null_congestion_pct\"),\n", + " (F.sum(F.when(F.col(\"Airport_fee\").isNull(), 1).otherwise(0)) / F.count(\"*\") * 100).alias(\"null_airport_fee_pct\")\n", + ")\n", + "\n", + "profile_pd = profile_df.toPandas().T\n", + "profile_pd.columns = [\"value\"]\n", + "profile_pd" + ] + }, + { + "cell_type": "markdown", + "id": "section_2_summary", + "metadata": {}, + "source": [ + "**Итоги профилирования:**\n", + "\n", + "| Проблема | Описание | Что будем делать |\n", + "|---|---|---|\n", + "| Отрицательные тарифы | Ошибки в данных (fare_amount < 0) | Фильтровать (F1) |\n", + "| Экстремальные суммы | Аномалии (total_amount > 1000) | Фильтровать (F2) |\n", + "| Даты вне диапазона | Записи за пределами 2024 года | Фильтровать (F3) |\n", + "| NULL в passenger_count | Не заполненные данные | Заполнять 0 и приводить к INT (N1) |\n", + "| NULL в RatecodeID | Не заполненные данные | Заполнять 99 (unknown) и приводить к INT (N2) |\n", + "| NULL в сборах | Не заполненные надбавки (surcharges/fees) | Заполнять 0.0 (N3, N4) |" + ] + }, + { + "cell_type": "markdown", + "id": "section_3_title", + "metadata": {}, + "source": [ + "## 3. Реестр правил трансформации\n", + "\n", + "Центральная идея модуля: фиксируем все правила явно перед кодом. Это аналог mapping specification в DWH. Если через месяц спросят «почему в silver нет строк с отрицательным fare?», ответ будет в реестре, а не зарыт в коде.\n", + "\n", + "| # | Колонка/область | Правило | Тип | Обоснование |\n", + "|---|---|---|---|---|\n", + "| **F1** | `fare_amount` | `fare_amount >= 0` | Фильтр | Отрицательные тарифы — явная аномалия |\n", + "| **F2** | `total_amount` | `total_amount <= 1000` | Фильтр | Экстремальные суммы — выбросы, искажающие среднее |\n", + "| **F3** | `tpep_pickup_datetime` | `[2024-01-01, 2025-01-01)` | Фильтр | Учебный диапазон данных — 2024 год |\n", + "| **N1** | `passenger_count` | `coalesce(cast(..., int), 0)` | Норм. | Приводим к целому, заменяем NULL на 0 |\n", + "| **N2** | `RatecodeID` | `coalesce(cast(..., int), 99)` | Норм. | Приводим к целому, заменяем NULL на 99 (Unknown) |\n", + "| **N3** | `congestion_surcharge` | `coalesce(..., 0.0)` | Норм. | Убираем NULL для корректности агрегатов |\n", + "| **N4** | `Airport_fee` | `coalesce(..., 0.0)` | Норм. | Убираем NULL для корректности агрегатов |\n", + "| **E1** | `trip_duration_minutes` | `(dropoff - pickup) / 60` | Обог. | Полезная вычисляемая метрика для анализа |\n", + "| **E2** | `pickup_zone/borough` | `LEFT JOIN zone_lookup` | Обог. | Понимаем, откуда поехала машина не по ID, а по имени |\n", + "| **E3** | `dropoff_zone/borough` | `LEFT JOIN zone_lookup` | Обог. | Понимаем, куда поехала машина не по ID, а по имени |" + ] + }, + { + "cell_type": "markdown", + "id": "section_4_title", + "metadata": {}, + "source": [ + "## 4. Применение правил — фильтрация\n", + "\n", + "Начинаем с фильтрации. Порядок осознанный: сначала фильтруем, потом нормализуем и обогащаем. Мы не тратим ресурсы на трансформацию строк, которые всё равно будут удалены." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_4_code", + "metadata": {}, + "outputs": [], + "source": [ + "bronze_df = spark.table(BRONZE_TABLE)\n", + "\n", + "filtered_df = (\n", + " bronze_df\n", + " .filter(F.col(\"fare_amount\") >= 0) # F1 \n", + " .filter(F.col(\"total_amount\") <= 1000) # F2 \n", + " .filter((F.col(\"tpep_pickup_datetime\") >= \"2024-01-01\") & (F.col(\"tpep_pickup_datetime\") < \"2025-01-01\")) # F3\n", + ")\n", + "\n", + "filtered_count = filtered_df.count()\n", + "removed_count = bronze_count - filtered_count\n", + "removed_pct = (removed_count / bronze_count) * 100\n", + "\n", + "print(f\"Строк до фильтрации: {bronze_count:,}\")\n", + "print(f\"Строк после фильтрации: {filtered_count:,}\")\n", + "print(f\"Удалено аномалий: {removed_count:,} ({removed_pct:.2f}%)\")" + ] + }, + { + "cell_type": "markdown", + "id": "section_4_note", + "metadata": {}, + "source": [ + "Обычно фильтры качества удаляют небольшую долю строк. Если фильтр убирает заметную часть данных — это сигнал перепроверить правило или источник." + ] + }, + { + "cell_type": "markdown", + "id": "section_5_title", + "metadata": {}, + "source": [ + "## 5. Применение правил — нормализация\n", + "\n", + "Приводим типы и обрабатываем NULL. В PostgreSQL мы бы использовали `ALTER COLUMN TYPE integer`, но в Lakehouse мы создаём новую версию данных между слоями." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_5_code", + "metadata": {}, + "outputs": [], + "source": [ + "normalized_df = (\n", + " filtered_df\n", + " .withColumn(\"passenger_count\", F.coalesce(F.col(\"passenger_count\").cast(\"int\"), F.lit(0))) # N1 \n", + " .withColumn(\"RatecodeID\", F.coalesce(F.col(\"RatecodeID\").cast(\"int\"), F.lit(99))) # N2 \n", + " .withColumn(\"congestion_surcharge\", F.coalesce(F.col(\"congestion_surcharge\"), F.lit(0.0))) # N3 \n", + " .withColumn(\"Airport_fee\", F.coalesce(F.col(\"Airport_fee\"), F.lit(0.0))) # N4\n", + ")\n", + "\n", + "normalized_df.select(\"passenger_count\", \"RatecodeID\", \"congestion_surcharge\", \"Airport_fee\").show(5)" + ] + }, + { + "cell_type": "markdown", + "id": "section_6_title", + "metadata": {}, + "source": [ + "## 6. Применение правил — обогащение\n", + "\n", + "Добавляем вычисляемые колонки и данные из lookup. Используем `LEFT JOIN`, чтобы не терять строки, если для какого-то ID не нашлось имени зоны. Это аналог денормализации в классическом DWH." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_6_code", + "metadata": {}, + "outputs": [], + "source": [ + "zone_lookup_df = spark.table(BRONZE_LOOKUP_TABLE)\n", + "\n", + "enriched_df = (\n", + " normalized_df\n", + " .withColumn(\"trip_duration_minutes\", (F.unix_timestamp(\"tpep_dropoff_datetime\") - F.unix_timestamp(\"tpep_pickup_datetime\")) / 60) # E1 \n", + " .join(zone_lookup_df.alias(\"pu_zone\"), F.col(\"PULocationID\") == F.col(\"pu_zone.LocationID\"), \"left\") # E2 \n", + " .withColumnRenamed(\"Borough\", \"pickup_borough\") \n", + " .withColumnRenamed(\"Zone\", \"pickup_zone\") \n", + " .drop(\"LocationID\", \"service_zone\") \n", + " .join(zone_lookup_df.alias(\"do_zone\"), F.col(\"DOLocationID\") == F.col(\"do_zone.LocationID\"), \"left\") # E3 \n", + " .withColumnRenamed(\"Borough\", \"dropoff_borough\") \n", + " .withColumnRenamed(\"Zone\", \"dropoff_zone\") \n", + " .drop(\"LocationID\", \"service_zone\")\n", + ")\n", + "\n", + "print(f\"Количество колонок после обогащения: {len(enriched_df.columns)}\")\n", + "enriched_df.select(\"tpep_pickup_datetime\", \"trip_duration_minutes\", \"pickup_zone\", \"dropoff_zone\").show(5, truncate=False)" + ] + }, + { + "cell_type": "markdown", + "id": "section_6_note", + "metadata": {}, + "source": [ + "**Обрати внимание:**\n", + "- **Отрицательная длительность:** если `tpep_dropoff_datetime` раньше `tpep_pickup_datetime`, колонка `trip_duration_minutes` будет отрицательной. Мы оставим это как наблюдение — это еще одна аномалия данных.\n", + "- **Смешанный регистр колонок:** оригинальные колонки (`VendorID`, `PULocationID`) сохраняют mixed case, тогда как наши новые колонки (`pickup_zone`, `trip_duration_minutes`) написаны в lowercase. `Iceberg` корректно обрабатывает такой набор колонок, но в реальных проектах лучше придерживаться единого стандарта." + ] + }, + { + "cell_type": "markdown", + "id": "section_7_title", + "metadata": {}, + "source": [ + "## 7. Создание silver-таблицы\n", + "\n", + "Используем `CREATE OR REPLACE TABLE ... AS SELECT` (CTAS) в namespace `lakehouse.silver`.\n", + "\n", + "**Важное предупреждение:** `CREATE OR REPLACE` — это полная перезапись таблицы. Все предыдущие данные и история snapshot-ов будут потеряны. В учебном silver это осознанный компромисс ради идемпотентности: повторный запуск ноутбука пересоздаёт таблицу с чистого листа. В Модуле 7 мы увидим, как делать изменения безопаснее.\n", + "\n", + "На 3 месяцах данных (~7 млн строк) эта операция может занять 2-4 минуты." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_7_code", + "metadata": {}, + "outputs": [], + "source": [ + "spark.sql(f\"CREATE NAMESPACE IF NOT EXISTS {SILVER_NAMESPACE}\")\n", + "\n", + "enriched_df.createOrReplaceTempView(\"silver_ready\")\n", + "\n", + "ctas_sql = f\"\"\"\n", + "CREATE OR REPLACE TABLE {SILVER_TABLE}\n", + "USING iceberg\n", + "AS\n", + "SELECT\n", + " VendorID,\n", + " tpep_pickup_datetime,\n", + " tpep_dropoff_datetime,\n", + " passenger_count,\n", + " trip_distance,\n", + " RatecodeID,\n", + " store_and_fwd_flag,\n", + " PULocationID,\n", + " DOLocationID,\n", + " payment_type,\n", + " fare_amount,\n", + " extra,\n", + " mta_tax,\n", + " tip_amount,\n", + " tolls_amount,\n", + " improvement_surcharge,\n", + " total_amount,\n", + " congestion_surcharge,\n", + " Airport_fee,\n", + " trip_duration_minutes,\n", + " pickup_borough,\n", + " pickup_zone,\n", + " dropoff_borough,\n", + " dropoff_zone\n", + "FROM silver_ready\n", + "\"\"\"\n", + "\n", + "print(f\"Выполняем CTAS для {SILVER_TABLE}...\")\n", + "spark.sql(ctas_sql)\n", + "\n", + "silver_df = spark.table(SILVER_TABLE)\n", + "print(f\"Silver-таблица готова. Строк: {silver_df.count():,}\")\n", + "silver_df.printSchema()" + ] + }, + { + "cell_type": "markdown", + "id": "section_8_title", + "metadata": {}, + "source": [ + "## 8. Проверки качества результата\n", + "\n", + "Ни одна трансформация не считается завершённой без проверок." + ] + }, + { + "cell_type": "markdown", + "id": "section_8a_title", + "metadata": {}, + "source": [ + "### 8a: Строковый уровень (Row-level)\n", + "Проверяем, что правила фильтрации и нормализации сработали для каждой строки." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_8a_code", + "metadata": {}, + "outputs": [], + "source": [ + "assert silver_df.filter(F.col(\"fare_amount\") < 0).count() == 0, \"Найдены отрицательные тарифы!\"\n", + "assert silver_df.filter(F.col(\"total_amount\") > 1000).count() == 0, \"Найдены экстремальные суммы!\"\n", + "assert silver_df.filter(F.col(\"passenger_count\").isNull()).count() == 0, \"Найдены NULL в passenger_count!\"\n", + "assert silver_df.filter(F.col(\"RatecodeID\").isNull()).count() == 0, \"Найдены NULL в RatecodeID!\"\n", + "\n", + "print(\"Row-level checks: OK\")" + ] + }, + { + "cell_type": "markdown", + "id": "section_8b_title", + "metadata": {}, + "source": [ + "### 8b: Схемный уровень (Schema-level)\n", + "Проверяем состав колонок и типы данных." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_8b_code", + "metadata": {}, + "outputs": [], + "source": [ + "expected_new_columns = [\"trip_duration_minutes\", \"pickup_borough\", \"pickup_zone\", \"dropoff_borough\", \"dropoff_zone\"]\n", + "for col in expected_new_columns:\n", + " assert col in silver_df.columns, f\"Колонка {col} отсутствует в silver!\"\n", + "\n", + "types = dict(silver_df.dtypes)\n", + "assert types[\"passenger_count\"] == \"int\", f\"Тип passenger_count должен быть int, а не {types['passenger_count']}\"\n", + "assert types[\"RatecodeID\"] == \"int\", f\"Тип RatecodeID должен быть int, а не {types['RatecodeID']}\"\n", + "\n", + "print(\"Schema-level checks: OK\")" + ] + }, + { + "cell_type": "markdown", + "id": "section_8c_title", + "metadata": {}, + "source": [ + "### 8c: Агрегатный уровень (Aggregate-level)\n", + "Сравниваем общие показатели bronze и silver, чтобы оценить масштаб изменений." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_8c_code", + "metadata": {}, + "outputs": [], + "source": [ + "comparison = spark.sql(f\"\"\"\n", + " WITH stats AS (\n", + " SELECT \n", + " 'bronze' as layer, \n", + " count(*) as rows, \n", + " avg(fare_amount) as avg_fare, \n", + " avg(trip_distance) as avg_dist\n", + " FROM {BRONZE_TABLE}\n", + " UNION ALL\n", + " SELECT \n", + " 'silver' as layer, \n", + " count(*) as rows, \n", + " avg(fare_amount) as avg_fare, \n", + " avg(trip_distance) as avg_dist\n", + " FROM {SILVER_TABLE}\n", + " )\n", + " SELECT \n", + " *,\n", + " CASE \n", + " WHEN layer = 'silver' THEN \n", + " (1 - rows * 1.0 / (SELECT rows FROM stats WHERE layer = 'bronze')) * 100 \n", + " ELSE 0 \n", + " END as filtered_pct\n", + " FROM stats\n", + "\"\"\").toPandas()\n", + "\n", + "comparison" + ] + }, + { + "cell_type": "markdown", + "id": "section_8c_note", + "metadata": {}, + "source": [ + "Ориентир: если фильтры удалили больше 5% строк, это повод пересмотреть правила или проверить данные." + ] + }, + { + "cell_type": "markdown", + "id": "section_9", + "metadata": {}, + "source": [ + "## 9. bronze vs silver — разделение ответственности\n", + "\n", + "Мы прошли путь от raw-файлов к очищенной silver-таблице. Почему бы не делать всё сразу при загрузке из raw?\n", + "\n", + "1. **Bronze хранит оригинал.** Если завтра выяснится, что правило фильтрации F1 было слишком жёстким, мы всегда можем пересоздать silver из bronze, не перечитывая исходные файлы из raw-зоны.\n", + "2. **Разделение сложности.** В bronze мы решаем задачу «прочитать и сохранить». В silver мы решаем задачу «очистить и обогатить». Это делает пайплайны проще в отладке.\n", + "3. **Единый источник правды.** Все downstream-потребители (аналитики, ML-инженеры) должны идти в silver. Bronze — это «внутренняя кухня» дата-инженера.\n", + "\n", + "| Свойство | Bronze | Silver |\n", + "|---|---|---|\n", + "| Состав данных | «Как в источнике» | Очищенные и обогащённые |\n", + "| NULL-значения | Сохранены | Обработаны (coalesce) |\n", + "| Типы данных | Технические (из raw) | Бизнес-ориентированные (INT для счетчиков) |\n", + "| Доп. колонки | Нет | Вычисляемые (длительность) и lookup (зоны) |\n", + "| Проверки качества | Минимальные (count) | Всесторонние (row/schema/agg) |\n", + "| Можно пересоздать из | Raw-файлов | Bronze-таблицы |" + ] + }, + { + "cell_type": "markdown", + "id": "section_10", + "metadata": {}, + "source": [ + "## 10. Самостоятельное задание\n", + "\n", + "Расширь silver-пайплайн тремя новыми правилами, пересоздай таблицу и проверь результат.\n", + "\n", + "**Фильтрация — F4:** удалить строки с `trip_distance = 0 AND total_amount != 0` (аномалия, найденная в Модуле 3).\n", + "\n", + "**Нормализация — N5:** привести `payment_type` из BIGINT в INT и заменить NULL на 0 (Unknown).\n", + "\n", + "**Обогащение — E4:** добавить вычисляемую колонку `speed_mph` — средняя скорость поездки в милях в час: `trip_distance / (trip_duration_minutes / 60)`. \n", + "*Подсказка: обработай деление на ноль — если trip_duration_minutes = 0, результат должен быть NULL, а не ошибка.*\n", + "\n", + "**Проверки:** после пересоздания silver напиши row-level assert для F4 и schema-level проверку для `speed_mph` и `payment_type`." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_10_f4", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: фильтрация F4\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_10_n5", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: нормализация N5\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_10_e4", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: обогащение E4\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_10_ctas", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: пересоздание silver-таблицы (CTAS)\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_10_checks", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: новые проверки качества\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_10_compare", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: сравнение количества строк до и после добавления F4\n" + ] + }, + { + "cell_type": "markdown", + "id": "section_11", + "metadata": {}, + "source": [ + "## 11. Checkpoint\n", + "\n", + "Проверь себя:\n", + "1. Чем silver отличается от bronze?\n", + "2. Зачем фиксировать правила перед кодом?\n", + "3. Какие три уровня проверок качества и зачем они нужны?\n", + "4. Почему мы использовали `LEFT JOIN`, а не `INNER JOIN` с lookup-таблицей?\n", + "5. Что делать, если фильтр внезапно удалил 50% строк?\n", + "6. Можно ли пересоздать silver, если через месяц мы обнаружим новую аномалию в данных?" + ] + }, + { + "cell_type": "markdown", + "id": "section_12", + "metadata": {}, + "source": [ + "## 12. Завершение\n", + "\n", + "Мы не удаляем silver-таблицу. В Модуле 6 мы будем читать её одновременно через Spark и Trino, а в Модуле 8 она понадобится для финальной практики." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "section_12_stop", + "metadata": {}, + "outputs": [], + "source": [ + "spark.stop()" + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": "Python 3 (ipykernel)", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.10.12" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} diff --git a/plans/README.md b/plans/README.md index 574145d..2ba2197 100644 --- a/plans/README.md +++ b/plans/README.md @@ -17,7 +17,7 @@ | 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 | +| 5. Слой silver и воспроизводимые трансформации | [module-05-silver-layer.md](./module-05-silver-layer.md) | Ready for validation | `notebooks/05_silver_layer.ipynb` | 2026-03-07 | | 6. Одна таблица, два движка: Spark и Trino | [module-06-spark-and-trino-on-same-table.md](./module-06-spark-and-trino-on-same-table.md) | Draft | `jupyter/Dockerfile`, `START_HERE.md`, `docs/stack_reference.md`, `notebooks/06_spark_and_trino_on_same_table.ipynb` | 2026-03-07 | | 7. Безопасная работа с таблицами: schema evolution и time travel | [module-07-safe-table-changes.md](./module-07-safe-table-changes.md) | Draft | `notebooks/07_safe_table_changes.ipynb` | 2026-03-07 | | 8. Базовое обслуживание таблиц и финальная практика | [module-08-maintenance-and-final-lab.md](./module-08-maintenance-and-final-lab.md) | Draft | `notebooks/08_maintenance_and_final_lab.ipynb` | 2026-03-07 | diff --git a/plans/module-05-silver-layer.md b/plans/module-05-silver-layer.md index c6ae736..f764d95 100644 --- a/plans/module-05-silver-layer.md +++ b/plans/module-05-silver-layer.md @@ -1,6 +1,6 @@ # Модуль 5. Слой silver и воспроизводимые трансформации -**Статус:** `Draft` +**Статус:** `Ready for validation` **Последнее обновление:** `2026-03-07` ## Цель