diff --git a/notebooks/04_bronze_with_iceberg.ipynb b/notebooks/04_bronze_with_iceberg.ipynb index 97291f2..8264559 100644 --- a/notebooks/04_bronze_with_iceberg.ipynb +++ b/notebooks/04_bronze_with_iceberg.ipynb @@ -564,21 +564,7 @@ "cell_type": "markdown", "id": "320c0199", "metadata": {}, - "source": [ - "## 8. Самостоятельное задание\n", - "\n", - "Используй raw-файл `taxi_zone_lookup.csv`, который уже лежит в `MinIO`.\n", - "\n", - "Что нужно сделать:\n", - "\n", - "1. Прочитать `taxi_zone_lookup.csv` из raw-зоны через `Spark`.\n", - "2. Создать таблицу `lakehouse.bronze.taxi_zone_lookup` с колонками `LocationID INT`, `Borough STRING`, `Zone STRING`, `service_zone STRING`.\n", - "3. Загрузить туда данные.\n", - "4. Проверить результат через `spark.table(...)`.\n", - "5. Посмотреть `snapshots` для этой таблицы.\n", - "\n", - "Подсказка: файл лежит по пути `s3a://lakehouse/raw/nyc_taxi/taxi_zone_lookup.csv`.\n" - ] + "source": "## 8. Вторая bronze-таблица: taxi_zone_lookup\n\nВ raw-зоне кроме Yellow Taxi лежит ещё один файл — `taxi_zone_lookup.csv`. Это справочник зон: каждому `LocationID` соответствует название района (`Borough`) и зоны (`Zone`).\n\nВ Модуле 5 мы будем обогащать silver-таблицу через `LEFT JOIN` с этим справочником. Сейчас загрузим его в bronze по тому же паттерну `CTAS`, но из CSV вместо Parquet." }, { "cell_type": "code", @@ -586,9 +572,7 @@ "id": "e94508bf", "metadata": {}, "outputs": [], - "source": [ - "# Ваш код: прочитай taxi_zone_lookup.csv из raw-зоны через Spark\n" - ] + "source": "raw_zone_df = (\n spark.read\n .option(\"header\", \"true\")\n .option(\"inferSchema\", \"true\")\n .csv(ZONE_LOOKUP_RAW_URI)\n)\n\nprint(f\"Строк в raw lookup: {raw_zone_df.count()}\")\nraw_zone_df.printSchema()\nraw_zone_df.show(5, truncate=False)" }, { "cell_type": "code", @@ -596,19 +580,15 @@ "id": "8e058aeb", "metadata": {}, "outputs": [], - "source": [ - "# Ваш код: создай таблицу lakehouse.bronze.taxi_zone_lookup\n" - ] + "source": "raw_zone_df.createOrReplaceTempView(\"raw_zone_lookup\")\n\nctas_lookup_sql = f\"\"\"\nCREATE OR REPLACE TABLE {BRONZE_LOOKUP_TABLE}\nUSING iceberg\nAS\nSELECT\n CAST(LocationID AS INT) AS LocationID,\n Borough,\n Zone,\n service_zone\nFROM raw_zone_lookup\n\"\"\"\n\nspark.sql(ctas_lookup_sql)\n\nlookup_df = spark.table(BRONZE_LOOKUP_TABLE)\nlookup_count = lookup_df.count()\n\nprint(f\"Строк в bronze lookup: {lookup_count}\")\nlookup_df.show(5, truncate=False)" }, { - "cell_type": "code", + "cell_type": "markdown", "execution_count": null, "id": "41c2c79a", "metadata": {}, "outputs": [], - "source": [ - "# Ваш код: загрузи данные в bronze lookup-таблицу\n" - ] + "source": "## 9. Самостоятельное задание\n\nПосмотри на внутреннюю структуру только что созданной `taxi_zone_lookup` — используй те же инструменты, что в Секциях 5 и 6.\n\n1. Выведи `snapshots` и `files` для `lakehouse.bronze.taxi_zone_lookup`.\n2. Сравни: сколько data files у lookup-таблицы и сколько у `nyc_taxi_yellow`? Почему такая разница?" }, { "cell_type": "code", @@ -616,9 +596,7 @@ "id": "80b392e1", "metadata": {}, "outputs": [], - "source": [ - "# Ваш код: проверь результат через spark.table(...)\n" - ] + "source": "# Ваш код: snapshots и files для taxi_zone_lookup\n" }, { "cell_type": "code", @@ -626,36 +604,19 @@ "id": "61e90959", "metadata": {}, "outputs": [], - "source": [ - "# Ваш код: посмотри snapshots для lakehouse.bronze.taxi_zone_lookup\n" - ] + "source": "# Ваш код: сравни количество data files в lookup и yellow taxi\n" }, { "cell_type": "markdown", "id": "24c04d60", "metadata": {}, - "source": [ - "## 9. Checkpoint\n", - "\n", - "Проверь себя:\n", - "\n", - "1. Чем отличается `spark.read.parquet(\"s3a://...\")` от `spark.table(\"lakehouse.bronze.nyc_taxi_yellow\")`?\n", - "2. Что хранит snapshot и зачем он нужен?\n", - "3. Где физически лежат данные таблицы и где лежат её метаданные?\n", - "4. Почему `bronze` это не то же самое, что `raw`?\n", - "5. Что произойдёт, если удалить один data file из `MinIO`, но metadata не поменять?\n", - "6. Можешь ли ты показать таблицу и через `Spark`, и через `MinIO Console`?\n" - ] + "source": "## 10. Checkpoint\n\nПроверь себя:\n\n1. Чем отличается `spark.read.parquet(\"s3a://...\")` от `spark.table(\"lakehouse.bronze.nyc_taxi_yellow\")`?\n2. Что хранит snapshot и зачем он нужен?\n3. Где физически лежат данные таблицы и где лежат её метаданные?\n4. Почему `bronze` это не то же самое, что `raw`?\n5. Что произойдёт, если удалить один data file из `MinIO`, но metadata не поменять?\n6. Можешь ли ты показать таблицу и через `Spark`, и через `MinIO Console`?" }, { "cell_type": "markdown", "id": "76d9061f", "metadata": {}, - "source": [ - "## 10. Завершение\n", - "\n", - "Мы намеренно не удаляем `lakehouse.bronze.nyc_taxi_yellow`. Эта таблица понадобится в Модуле 5, где мы будем строить `silver`-слой.\n" - ] + "source": "## 11. Завершение\n\nМы намеренно не удаляем bronze-таблицы (`nyc_taxi_yellow` и `taxi_zone_lookup`). Они понадобятся в Модуле 5, где мы будем строить silver-слой." }, { "cell_type": "code", diff --git a/notebooks/05_silver_layer.ipynb b/notebooks/05_silver_layer.ipynb index 03f6d60..6f26752 100644 --- a/notebooks/05_silver_layer.ipynb +++ b/notebooks/05_silver_layer.ipynb @@ -77,21 +77,7 @@ "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.\")" - ] + "source": "assert spark.catalog.tableExists(BRONZE_TABLE), (\n f\"Таблица {BRONZE_TABLE} не найдена. Вернись в Модуль 4 и выполни загрузку в bronze.\"\n)\n\nassert spark.catalog.tableExists(BRONZE_LOOKUP_TABLE), (\n f\"Таблица {BRONZE_LOOKUP_TABLE} не найдена. \"\n \"Вернись в Модуль 4 и выполни Секцию 8.\"\n)\n\nbronze_count = spark.table(BRONZE_TABLE).count()\nassert bronze_count > 0, f\"Таблица {BRONZE_TABLE} пуста.\"\n\nprint(f\"Обе bronze-таблицы на месте. В основной таблице {bronze_count:,} строк. Теперь строим silver.\")" }, { "cell_type": "markdown", @@ -624,4 +610,4 @@ }, "nbformat": 4, "nbformat_minor": 5 -} +} \ No newline at end of file diff --git a/plans/module-04-bronze-with-iceberg.md b/plans/module-04-bronze-with-iceberg.md index 8d16f8d..d8b92cf 100644 --- a/plans/module-04-bronze-with-iceberg.md +++ b/plans/module-04-bronze-with-iceberg.md @@ -70,9 +70,9 @@ Bronze нужен Модулю 5. Spark-сессия останавливается. Явное пояснение в финале ноутбука: «Мы намеренно не удаляем bronze-таблицы. В Модуле 5 мы будем строить silver на основе bronze.» -### Taxi zone lookup: самостоятельное задание +### Taxi zone lookup: демо-секция -Студент создаёт `lakehouse.bronze.taxi_zone_lookup` как упражнение. Это подготовка к join в Модуле 5 (silver). +Создание `lakehouse.bronze.taxi_zone_lookup` входит в демонстрационный код ноутбука (Секция 8), а не в самостоятельное задание. Причина: Модуль 5 использует lookup для JOIN в демо-коде — если lookup не создан, студент не может пройти Модуль 5, даже выполняя только готовые ячейки. Принцип: демо-ячейки последующих модулей не должны зависеть от домашних заданий предыдущих. ## План работ @@ -160,11 +160,16 @@ Bronze нужен Модулю 5. Spark-сессия останавливает - Bronze ≠ raw. Raw — неизменяемый архив файлов. Bronze — управляемая таблица, которую можно пересоздать из raw. - В Модуле 5 — silver: очистка, трансформации, бизнес-логика. -### Секция 7: Самостоятельное задание -- **[md]** Инструкция: (1) прочитать `taxi_zone_lookup.csv` из raw-зоны, (2) создать `lakehouse.bronze.taxi_zone_lookup` с DDL (LocationID INT, Borough STRING, Zone STRING, service_zone STRING), (3) загрузить данные, (4) проверить через `spark.table(...)`, (5) осмотреть snapshots. -- **[code]** 5 пустых ячеек `# Ваш код: ...`. +### Секция 7: Вторая bronze-таблица: taxi_zone_lookup +- **[md]** В raw-зоне кроме Yellow Taxi лежит `taxi_zone_lookup.csv` — справочник зон. В Модуле 5 он понадобится для JOIN. Загружаем в bronze по тому же паттерну CTAS, но из CSV. +- **[code]** Чтение CSV через `spark.read.option("header", "true").option("inferSchema", "true").csv(...)`. Вывод count, schema, preview. +- **[code]** CTAS: `CREATE OR REPLACE TABLE lakehouse.bronze.taxi_zone_lookup USING iceberg AS SELECT CAST(LocationID AS INT) AS LocationID, Borough, Zone, service_zone FROM raw_zone_lookup`. Верификация: count, show. -### Секция 8: Checkpoint +### Секция 8: Самостоятельное задание +- **[md]** Инструкция: осмотреть внутреннюю структуру `taxi_zone_lookup` теми же инструментами, что в Секциях 5 и 6. (1) Вывести snapshots и files. (2) Сравнить количество data files с `nyc_taxi_yellow` и объяснить разницу. +- **[code]** 2 пустых ячейки `# Ваш код: ...`. + +### Секция 9: Checkpoint - **[md]** Вопросы для самопроверки: 1. Чем отличается `spark.read.parquet("s3a://...")` от `spark.table("lakehouse.bronze.nyc_taxi_yellow")`? 2. Что хранит snapshot и зачем он нужен? @@ -173,8 +178,8 @@ Bronze нужен Модулю 5. Spark-сессия останавливает 5. Что произойдёт, если удалить один data file из MinIO, но не изменить metadata? 6. Можете ли вы показать таблицу в MinIO Console (через браузер)? -### Секция 9: Завершение -- **[md]** «Мы намеренно не удаляем bronze-таблицы. В Модуле 5 мы будем строить silver-слой на основе `lakehouse.bronze.nyc_taxi_yellow`.» +### Секция 10: Завершение +- **[md]** «Мы намеренно не удаляем bronze-таблицы (`nyc_taxi_yellow` и `taxi_zone_lookup`). Они понадобятся в Модуле 5, где мы будем строить silver-слой.» - **[code]** `spark.stop()`. ### Дизайн-решения по ноутбуку diff --git a/plans/module-05-silver-layer.md b/plans/module-05-silver-layer.md index f764d95..225307f 100644 --- a/plans/module-05-silver-layer.md +++ b/plans/module-05-silver-layer.md @@ -88,9 +88,9 @@ LEFT JOIN (не INNER), чтобы не терять строки с LocationID - **Схемный (schema-level):** наличие новых колонок, проверка типов (passenger_count = INT, RatecodeID = INT) - **Агрегатный (aggregate-level):** таблица-сравнение bronze vs silver (row count, % отфильтрованных, средние fare/distance) -### Taxi zone lookup: жёсткий prerequisite +### Taxi zone lookup: prerequisite из демо-секции Модуля 4 -Lookup — самостоятельное задание Модуля 4. Модуль 5 работает только с bronze-таблицами и не должен лезть в raw-слой — это нарушило бы принцип разделения ответственности между слоями (PRD Risk #4). Стратегия: жёсткий assert при отсутствии `lakehouse.bronze.taxi_zone_lookup` с понятным сообщением и отсылкой к Модулю 4: «Таблица `lakehouse.bronze.taxi_zone_lookup` не найдена. Вернись в Модуль 4 и выполни самостоятельное задание (Секция 8).» +Lookup создаётся в демонстрационном коде Модуля 4 (Секция 8), а не в самостоятельном задании. Это гарантирует, что студент, прогнавший готовые ячейки Модуля 4, имеет обе bronze-таблицы. Модуль 5 работает только с bronze-таблицами и не должен лезть в raw-слой — это нарушило бы принцип разделения ответственности между слоями (PRD Risk #4). Стратегия: жёсткий assert при отсутствии `lakehouse.bronze.taxi_zone_lookup` с понятным сообщением и отсылкой к Модулю 4: «Таблица `lakehouse.bronze.taxi_zone_lookup` не найдена. Вернись в Модуль 4 и выполни Секцию 8.» ### Cleanup: нет @@ -119,7 +119,7 @@ Silver нужен Модулям 6 и 8. Spark-сессия останавлив ### Секция 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). +- **[code]** Assert: `spark.table(BRONZE_TABLE).count() > 0` с отсылкой к Модулю 4. Assert: `spark.table(BRONZE_LOOKUP_TABLE)` существует — при отсутствии жёсткая ошибка с отсылкой к Секции 8 Модуля 4. - **[md]** Обе bronze-таблицы на месте. Теперь строим silver. ### Секция 2: Профиль bronze — что нужно трансформировать @@ -272,7 +272,7 @@ speed_mph DOUBLE (E4: новая, trip_distance / (tr ## Риски -- **Taxi zone lookup отсутствует.** Самостоятельное задание Модуля 4 может быть не выполнено. Жёсткий assert с отсылкой к Модулю 4. Модуль 5 не создаёт bronze-таблицы самостоятельно — это нарушило бы разделение ответственности между слоями. +- **Taxi zone lookup отсутствует.** Lookup создаётся в демо-секции Модуля 4 (Секция 8), поэтому риск минимален — студент должен был прогнать готовые ячейки. Жёсткий 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 — длительность отрицательна. Упомянуть как наблюдение, не фильтровать в демо (потенциально — для самостоятельного задания или будущего правила).