From 6567a44b74f6f99a60e725b70e92cf7d6df74432 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Mar 2026 19:22:04 +0300 Subject: [PATCH] =?UTF-8?q?feat(course):=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2?= =?UTF-8?q?=D0=BB=D0=B5=D0=BD=20=D0=BD=D0=BE=D1=83=D1=82=D0=B1=D1=83=D0=BA?= =?UTF-8?q?=20=D0=BC=D0=BE=D0=B4=D1=83=D0=BB=D1=8F=20bronze=20=D1=81=20Ice?= =?UTF-8?q?berg?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - нужен практический модуль 4, который переводит студента от raw-файлов к первой управляемой Iceberg-таблице. - Что: - добавлен ноутбук notebooks/04_bronze_with_iceberg.ipynb с CTAS в bronze, разбором metadata tables и структурой хранения в MinIO. - добавлены пояснения про namespace, сравнение raw vs Iceberg и разбор snapshot-полей по итогам ревью. - обновлен plans/README.md строкой для модуля 4. - Проверка: - синтаксис всех code-ячеек ноутбука проверен локальной компиляцией. - тесты пройдены пользователем. --- notebooks/04_bronze_with_iceberg.ipynb | 692 +++++++++++++++++++++++++ plans/README.md | 1 + 2 files changed, 693 insertions(+) create mode 100644 notebooks/04_bronze_with_iceberg.ipynb diff --git a/notebooks/04_bronze_with_iceberg.ipynb b/notebooks/04_bronze_with_iceberg.ipynb new file mode 100644 index 0000000..97291f2 --- /dev/null +++ b/notebooks/04_bronze_with_iceberg.ipynb @@ -0,0 +1,692 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "id": "2a1c9579", + "metadata": {}, + "source": [ + "# Модуль 4. Первая рабочая Iceberg-таблица и слой bronze\n", + "\n", + "В этом модуле мы переходим от raw-файлов в `MinIO` к первой управляемой таблице `Iceberg`.\n", + "\n", + "Что делаем на практике:\n", + "\n", + "- проверяем, что raw-данные из Модуля 3 доступны в `MinIO`;\n", + "- создаём namespace `lakehouse.bronze`;\n", + "- загружаем raw `Yellow Taxi` в `Iceberg`-таблицу через `CTAS`;\n", + "- сравниваем чтение по физическому пути и чтение через каталог;\n", + "- смотрим `snapshots`, `history`, `files` и физическую структуру таблицы в `MinIO`.\n", + "\n", + "После этого ноутбука ты должен понимать, почему `bronze`-слой это уже не «просто папка с parquet», а управляемая таблица с каталогом и metadata.\n" + ] + }, + { + "cell_type": "markdown", + "id": "e13e7b2b", + "metadata": {}, + "source": "## 0. Перед стартом\n\nПеред выполнением ноутбука:\n\n- подними стенд по `START_HERE.md`;\n- пройди `notebooks/03_raw_ingest_and_first_read.ipynb`;\n- убедись, что raw-файлы уже лежат в `s3a://lakehouse/raw/nyc_taxi/`.\n\nСвязка с привычным `PostgreSQL / Greenplum`:\n\n| | Классический DWH (`Greenplum` / `PostgreSQL`) | Lakehouse |\n| --- | --- | --- |\n| Raw-данные | Загружены в `stg`-таблицу или внешний staging | Файлы в object storage (`raw/`) |\n| Первый управляемый слой | Явный DDL + `INSERT INTO ... SELECT` из `stg` | `CREATE TABLE ... USING iceberg AS SELECT` из raw (CTAS) |\n| Метаданные | `pg_catalog` внутри СУБД | JDBC catalog + metadata files в `MinIO` |\n| История изменений | Нет встроенной истории таблицы | `snapshots`, `manifests`, metadata JSON |\n\nЗдесь важно различать два уровня:\n\n- raw: исходные файлы, которые можно перечитать заново;\n- bronze: первая управляемая таблица, которую удобно читать из `Spark` и позже из `Trino`." + }, + { + "cell_type": "code", + "execution_count": null, + "id": "33cd36c0", + "metadata": {}, + "outputs": [], + "source": [ + "from pathlib import PurePosixPath\n", + "\n", + "import boto3\n", + "import pandas as pd\n", + "from pyspark.sql import SparkSession, functions as F\n", + "\n", + "BUCKET_NAME = \"lakehouse\"\n", + "RAW_PREFIX = \"raw/nyc_taxi/\"\n", + "RAW_URI = f\"s3a://{BUCKET_NAME}/{RAW_PREFIX}\"\n", + "RAW_YELLOW_URI = f\"{RAW_URI}yellow_tripdata_*.parquet\"\n", + "ZONE_LOOKUP_RAW_URI = f\"{RAW_URI}taxi_zone_lookup.csv\"\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", + "spark = SparkSession.builder .appName(\"module-04-bronze-with-iceberg\") .getOrCreate()\n", + "\n", + "spark.sparkContext.setLogLevel(\"ERROR\")\n", + "\n", + "s3 = boto3.client(\n", + " \"s3\",\n", + " endpoint_url=\"http://minio:9000\",\n", + " aws_access_key_id=\"minioadmin\",\n", + " aws_secret_access_key=\"minioadmin\",\n", + ")\n", + "\n", + "\n", + "def format_bytes(num_bytes: int) -> str:\n", + " units = [\"B\", \"KB\", \"MB\", \"GB\", \"TB\"]\n", + " value = float(num_bytes)\n", + " for unit in units:\n", + " if value < 1024 or unit == units[-1]:\n", + " return f\"{value:.1f} {unit}\"\n", + " value /= 1024\n", + "\n", + "\n", + "def list_objects(bucket: str, prefix: str) -> list[dict]:\n", + " paginator = s3.get_paginator(\"list_objects_v2\")\n", + " objects = []\n", + " for page in paginator.paginate(Bucket=bucket, Prefix=prefix):\n", + " objects.extend(page.get(\"Contents\", []))\n", + " return objects\n", + "\n", + "\n", + "def get_table_location(table_name: str) -> str:\n", + " details = spark.sql(f\"DESCRIBE TABLE EXTENDED {table_name}\").toPandas()\n", + " location_row = details.loc[details[\"col_name\"] == \"Location\", \"data_type\"]\n", + " assert not location_row.empty, f\"Не удалось найти Location для {table_name}\"\n", + " return location_row.iloc[0]\n", + "\n", + "\n", + "print(\"Spark и MinIO-клиент готовы\")\n", + "print(f\"RAW URI: {RAW_URI}\")\n", + "print(f\"Bronze table: {BRONZE_TABLE}\")\n" + ] + }, + { + "cell_type": "markdown", + "id": "e195a355", + "metadata": {}, + "source": "## 1. Проверка raw перед загрузкой в bronze\n\nСначала убеждаемся, что raw-данные из Модуля 3 действительно доступны в `MinIO`. Если их нет, строить `bronze` просто не из чего.\n\nКлючевое различие:\n\n- `raw` читается по физическому пути;\n- `bronze` будет читаться уже как логическая таблица через каталог." + }, + { + "cell_type": "code", + "execution_count": null, + "id": "9e69d44a", + "metadata": {}, + "outputs": [], + "source": [ + "raw_objects = list_objects(BUCKET_NAME, RAW_PREFIX)\n", + "assert raw_objects, (\n", + " \"Raw-зона пуста. Сначала пройди Модуль 3 и загрузи data bundle в MinIO.\"\n", + ")\n", + "\n", + "raw_report = pd.DataFrame(\n", + " [\n", + " {\n", + " \"key\": obj[\"Key\"],\n", + " \"size_bytes\": obj[\"Size\"],\n", + " \"size_human\": format_bytes(obj[\"Size\"]),\n", + " }\n", + " for obj in raw_objects\n", + " ]\n", + ").sort_values(\"key\").reset_index(drop=True)\n", + "\n", + "yellow_object_count = int(raw_report[\"key\"].str.contains(r\"yellow_tripdata_.*\\.parquet\", regex=True).sum())\n", + "assert yellow_object_count > 0, (\n", + " \"В raw-зоне не найдены yellow_tripdata_*.parquet. Сначала пройди Модуль 3.\"\n", + ")\n", + "assert (raw_report[\"key\"] == f\"{RAW_PREFIX}taxi_zone_lookup.csv\").any(), (\n", + " \"В raw-зоне не найден taxi_zone_lookup.csv. Он нужен для задания этого модуля и join в Модуле 5.\"\n", + ")\n", + "\n", + "print(f\"Всего объектов в raw-зоне: {len(raw_report)}\")\n", + "print(f\"Parquet-файлов Yellow Taxi: {yellow_object_count}\")\n", + "raw_report\n" + ] + }, + { + "cell_type": "markdown", + "id": "f97b8ec8", + "metadata": {}, + "source": "## 2. Чтение raw-файлов через Spark\n\nСейчас мы ещё работаем с raw напрямую: `spark.read.parquet(\"s3a://...\")`.\n\nЭто не таблица `Iceberg`, а просто набор parquet-файлов по физическому пути. Нам нужно посмотреть на схему и подготовить явный список колонок для `CTAS`, чтобы не полагаться на `SELECT *`." + }, + { + "cell_type": "code", + "execution_count": null, + "id": "61f18800", + "metadata": {}, + "outputs": [], + "source": [ + "raw_yellow_df = spark.read.parquet(RAW_YELLOW_URI)\n", + "raw_row_count = raw_yellow_df.count()\n", + "\n", + "print(f\"Прочитано raw-строк: {raw_row_count:,}\")\n", + "print(f\"Количество колонок: {len(raw_yellow_df.columns)}\")\n", + "raw_yellow_df.show(10, truncate=False)\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "a04c71e6", + "metadata": {}, + "outputs": [], + "source": [ + "raw_yellow_df.printSchema()\n", + "\n", + "pd.DataFrame(raw_yellow_df.dtypes, columns=[\"column\", \"spark_type\"])\n" + ] + }, + { + "cell_type": "markdown", + "id": "ce148bf6", + "metadata": {}, + "source": "Ниже список колонок, который мы явно зафиксируем в `CTAS`.\n\nПочему это лучше, чем `SELECT *`:\n\n- порядок и состав колонок становятся явными;\n- проще заметить неожиданные изменения источника;\n- загрузка в `bronze` остаётся сознательным решением, а не непрозрачным автоматизмом." + }, + { + "cell_type": "code", + "execution_count": null, + "id": "305ec274", + "metadata": {}, + "outputs": [], + "source": [ + "expected_columns = [\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", + "]\n", + "\n", + "missing_columns = sorted(set(expected_columns).difference(raw_yellow_df.columns))\n", + "assert not missing_columns, (\n", + " \"В raw-схеме не хватает ожидаемых колонок: \"\n", + " f\"{missing_columns}. Проверь data bundle и инструкцию из START_HERE.md.\"\n", + ")\n", + "\n", + "pd.DataFrame({\n", + " \"column\": expected_columns,\n", + " \"present_in_raw\": [column in raw_yellow_df.columns for column in expected_columns],\n", + "})\n" + ] + }, + { + "cell_type": "markdown", + "id": "350d2087", + "metadata": {}, + "source": "## 3. Создание `bronze` namespace и `Iceberg`-таблицы\n\nСначала создаём namespace `lakehouse.bronze`, а затем выполняем `CTAS` (`CREATE OR REPLACE TABLE ... USING iceberg AS SELECT ...`).\n\nЧто такое namespace в `Iceberg`:\n\n- это логический контейнер для таблиц;\n- если ты работал с `PostgreSQL` — это ближайший аналог `schema` внутри базы;\n- запись `lakehouse.bronze.nyc_taxi_yellow` читается как `catalog.namespace.table`.\n\nПочему здесь используем именно `bronze`, а не `module_04`:\n\n- имя сразу отражает роль слоя в пайплайне `raw -> bronze -> silver`;\n- Модуль 5 будет читать именно `lakehouse.bronze.*`, и семантическое имя делает маршрут прозрачным;\n- в отличие от демонстрационного `module_02`, этот слой не временный и будет использоваться в следующих модулях.\n\nЭто удобный учебный компромисс:\n\n- один шаг создаёт таблицу и сразу наполняет её;\n- повторный запуск секции остаётся идемпотентным;\n- минус: `CREATE OR REPLACE` пересоздаёт таблицу и сбрасывает старую историю snapshot-ов.\n\nДля первого `bronze`-слоя это допустимо. В Модуле 7 мы отдельно разберём, почему в рабочих сценариях нужно аккуратнее относиться к истории таблицы." + }, + { + "cell_type": "code", + "execution_count": null, + "id": "a52a03b4", + "metadata": {}, + "outputs": [], + "source": [ + "spark.sql(f\"CREATE NAMESPACE IF NOT EXISTS {BRONZE_NAMESPACE}\")\n", + "raw_yellow_df.createOrReplaceTempView(\"raw_yellow\")\n", + "\n", + "ctas_sql = f\"\"\"\n", + "CREATE OR REPLACE TABLE {BRONZE_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", + "FROM raw_yellow\n", + "\"\"\"\n", + "\n", + "print(\"Temp view ready: raw_yellow\")\n", + "print(ctas_sql)\n" + ] + }, + { + "cell_type": "markdown", + "id": "773150ad", + "metadata": {}, + "source": [ + "Справочная форма DDL, более похожая на привычный `PostgreSQL`-стиль, могла бы выглядеть так:\n", + "\n", + "```sql\n", + "CREATE TABLE lakehouse.bronze.nyc_taxi_yellow (\n", + " VendorID BIGINT,\n", + " tpep_pickup_datetime TIMESTAMP,\n", + " tpep_dropoff_datetime TIMESTAMP,\n", + " passenger_count DOUBLE,\n", + " trip_distance DOUBLE,\n", + " RatecodeID DOUBLE,\n", + " store_and_fwd_flag STRING,\n", + " PULocationID BIGINT,\n", + " DOLocationID BIGINT,\n", + " payment_type BIGINT,\n", + " fare_amount DOUBLE,\n", + " extra DOUBLE,\n", + " mta_tax DOUBLE,\n", + " tip_amount DOUBLE,\n", + " tolls_amount DOUBLE,\n", + " improvement_surcharge DOUBLE,\n", + " total_amount DOUBLE,\n", + " congestion_surcharge DOUBLE,\n", + " Airport_fee DOUBLE\n", + ") USING iceberg;\n", + "```\n", + "\n", + "Но в этом модуле мы используем именно `CTAS`, потому что хотим одной операцией и создать таблицу, и загрузить в неё raw-срез.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "3fa6f522", + "metadata": {}, + "outputs": [], + "source": [ + "# Для полного набора из нескольких месяцев эта ячейка может выполняться 1-3 минуты.\n", + "spark.sql(ctas_sql)\n", + "\n", + "bronze_df = spark.table(BRONZE_TABLE)\n", + "bronze_row_count = bronze_df.count()\n", + "\n", + "print(f\"Строк в raw: {raw_row_count:,}\")\n", + "print(f\"Строк в bronze: {bronze_row_count:,}\")\n", + "assert bronze_row_count == raw_row_count, \"Количество строк в bronze не совпало с raw\"\n" + ] + }, + { + "cell_type": "markdown", + "id": "70dead5a", + "metadata": {}, + "source": "## 4. Читаем уже не raw, а управляемую таблицу\n\nТеперь ключевое отличие:\n\n- `spark.read.parquet(\"s3a://...\")` читает набор файлов по физическому пути;\n- `spark.table(\"lakehouse.bronze.nyc_taxi_yellow\")` читает логическую таблицу через каталог `Iceberg`.\n\nПо результату оба варианта могут вернуть похожие строки. Но второй вариант работает с каталогом и metadata: знает схему, историю изменений и согласованное состояние таблицы." + }, + { + "cell_type": "code", + "execution_count": null, + "id": "252cba12", + "metadata": {}, + "outputs": [], + "source": [ + "bronze_df.show(10, truncate=False)\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "abfd4ee4", + "metadata": {}, + "outputs": [], + "source": [ + "comparison_df = pd.DataFrame(\n", + " [\n", + " {\n", + " \"mode\": \"raw path\",\n", + " \"reader\": f\"spark.read.parquet('{RAW_YELLOW_URI}')\",\n", + " \"row_count\": raw_row_count,\n", + " \"object\": \"набор parquet-файлов\",\n", + " },\n", + " {\n", + " \"mode\": \"catalog table\",\n", + " \"reader\": f\"spark.table('{BRONZE_TABLE}')\",\n", + " \"row_count\": bronze_row_count,\n", + " \"object\": \"Iceberg-таблица\",\n", + " },\n", + " ]\n", + ")\n", + "\n", + "comparison_df\n" + ] + }, + { + "cell_type": "markdown", + "id": "5b4dffa6", + "metadata": {}, + "source": "Эту разницу полезно зафиксировать в одной таблице:\n\n| Свойство | Raw-файлы | Iceberg-таблица |\n| --- | --- | --- |\n| Каталогизация | Нет | Да (`PostgreSQL` JDBC catalog) |\n| История изменений | Нет | `Snapshots` |\n| Контроль схемы | Выводится из файлов | Зафиксирована в metadata |\n| Чтение из `Trino` | Невозможно без внешней таблицы | Да |\n| Атомарность записи | Нет | Да |" + }, + { + "cell_type": "markdown", + "id": "ccb4fb69", + "metadata": {}, + "source": [ + "## 5. Что знает `Iceberg` о таблице: snapshots, history, files\n", + "\n", + "У `Iceberg` есть встроенные metadata-таблицы. Это один из самых наглядных ответов на вопрос, почему таблица не сводится к «папке с parquet-файлами».\n", + "\n", + "Нас интересуют три представления:\n", + "\n", + "- `.snapshots` — какие snapshots есть у таблицы;\n", + "- `.history` — как менялось текущее состояние таблицы;\n", + "- `.files` — какие data files входят в текущий snapshot.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "03fa2fc4", + "metadata": {}, + "outputs": [], + "source": [ + "spark.sql(\n", + " f\"SELECT committed_at, snapshot_id, parent_id, operation, manifest_list, summary \"\n", + " f\"FROM {BRONZE_TABLE}.snapshots ORDER BY committed_at DESC\"\n", + ").toPandas()\n" + ] + }, + { + "cell_type": "markdown", + "id": "3f048c15", + "metadata": {}, + "source": "Как читать поля в `.snapshots`:\n\n- `committed_at` — когда snapshot стал текущим;\n- `snapshot_id` — уникальный идентификатор версии таблицы;\n- `operation` — какая операция создала snapshot, например `append` или `replace`;\n- `summary` — краткое описание результата операции: сколько файлов добавлено, сколько строк записано и т.д.\n\nНа практике это можно воспринимать как именованный checkpoint таблицы: каждая запись в истории фиксирует согласованное состояние, к которому движок может обратиться." + }, + { + "cell_type": "code", + "execution_count": null, + "id": "db404328", + "metadata": {}, + "outputs": [], + "source": [ + "spark.sql(\n", + " f\"SELECT made_current_at, snapshot_id, parent_id, is_current_ancestor \"\n", + " f\"FROM {BRONZE_TABLE}.history ORDER BY made_current_at DESC\"\n", + ").toPandas()\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "dc0f3c91", + "metadata": {}, + "outputs": [], + "source": [ + "files_df = spark.sql(\n", + " f\"\"\"\n", + " SELECT\n", + " content,\n", + " file_path,\n", + " file_format,\n", + " record_count,\n", + " file_size_in_bytes\n", + " FROM {BRONZE_TABLE}.files\n", + " ORDER BY file_path\n", + " \"\"\"\n", + ")\n", + "\n", + "files_pd = files_df.toPandas()\n", + "files_pd[\"file_size_human\"] = files_pd[\"file_size_in_bytes\"].map(format_bytes)\n", + "files_pd\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "a9cfcaa6", + "metadata": {}, + "outputs": [], + "source": [ + "file_summary = pd.DataFrame(\n", + " [\n", + " {\n", + " \"metric\": \"data_files_count\",\n", + " \"value\": len(files_pd),\n", + " },\n", + " {\n", + " \"metric\": \"total_records_in_files\",\n", + " \"value\": int(files_pd[\"record_count\"].sum()) if not files_pd.empty else 0,\n", + " },\n", + " {\n", + " \"metric\": \"total_file_size_bytes\",\n", + " \"value\": int(files_pd[\"file_size_in_bytes\"].sum()) if not files_pd.empty else 0,\n", + " },\n", + " {\n", + " \"metric\": \"total_file_size_human\",\n", + " \"value\": format_bytes(int(files_pd[\"file_size_in_bytes\"].sum())) if not files_pd.empty else \"0 B\",\n", + " },\n", + " ]\n", + ")\n", + "\n", + "file_summary\n" + ] + }, + { + "cell_type": "markdown", + "id": "bc85ca8f", + "metadata": {}, + "source": "`Iceberg` хранит не только путь к данным, но и состояние таблицы как набора metadata-артефактов.\n\nПрактический смысл:\n\n- snapshot фиксирует согласованную версию таблицы;\n- manifests описывают, какие data files входят в snapshot;\n- каталог указывает на актуальный metadata-файл таблицы.\n\nИменно поэтому удалить один parquet-файл «руками» из `MinIO` опаснее, чем кажется: metadata продолжит ссылаться на него, а таблица станет неконсистентной." + }, + { + "cell_type": "markdown", + "id": "d48ea169", + "metadata": {}, + "source": [ + "## 6. Физическая структура таблицы в `MinIO`\n", + "\n", + "Теперь посмотрим на ту же таблицу с физической стороны. Нас интересуют две группы объектов:\n", + "\n", + "- `data/` — parquet-файлы с данными;\n", + "- `metadata/` — metadata JSON, manifest list и manifests.\n", + "\n", + "Это хороший момент, чтобы связать логическую таблицу и физическую структуру в object storage.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "9d2d088d", + "metadata": {}, + "outputs": [], + "source": [ + "bronze_location = get_table_location(BRONZE_TABLE)\n", + "bronze_prefix = bronze_location.replace(f\"s3a://{BUCKET_NAME}/\", \"\").rstrip(\"/\") + \"/\"\n", + "bronze_objects = list_objects(BUCKET_NAME, bronze_prefix)\n", + "assert bronze_objects, f\"По пути {bronze_location} не найдены объекты таблицы\"\n", + "\n", + "bronze_storage_report = pd.DataFrame(\n", + " [\n", + " {\n", + " \"key\": obj[\"Key\"],\n", + " \"kind\": \"DATA\" if \"/data/\" in obj[\"Key\"] else \"META\" if \"/metadata/\" in obj[\"Key\"] else \"OTHER\",\n", + " \"size_bytes\": obj[\"Size\"],\n", + " \"size_human\": format_bytes(obj[\"Size\"]),\n", + " }\n", + " for obj in bronze_objects\n", + " ]\n", + ").sort_values([\"kind\", \"key\"]).reset_index(drop=True)\n", + "\n", + "print(f\"Table location: {bronze_location}\")\n", + "bronze_storage_report\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "c2fbc7e1", + "metadata": {}, + "outputs": [], + "source": [ + "storage_summary = bronze_storage_report.groupby(\"kind\", as_index=False)[\"size_bytes\"].sum()\n", + "storage_summary[\"size_human\"] = storage_summary[\"size_bytes\"].map(format_bytes)\n", + "storage_summary\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "2fb4911d", + "metadata": {}, + "outputs": [], + "source": [ + "root = PurePosixPath(bronze_prefix.rstrip(\"/\"))\n", + "rendered_lines = [str(root)]\n", + "\n", + "for key in bronze_storage_report[\"key\"]:\n", + " rel = PurePosixPath(key).relative_to(root)\n", + " parts = rel.parts\n", + " for depth, part in enumerate(parts):\n", + " prefix = \" \" * depth + (\"- \" if depth else \"\")\n", + " rendered = prefix + part\n", + " if rendered not in rendered_lines:\n", + " rendered_lines.append(rendered)\n", + "\n", + "print(\"Структура каталога таблицы:\\n\")\n", + "print(\"\\n\".join(rendered_lines))\n" + ] + }, + { + "cell_type": "markdown", + "id": "ef5b24bf", + "metadata": {}, + "source": "Обычно после первого `CTAS` ты увидишь примерно такую картину:\n\n```text\nwarehouse/bronze/nyc_taxi_yellow/\n - data\n - *.parquet\n - metadata\n - v1.metadata.json\n - snap-*.avro\n - *.avro\n```\n\nСмысл структуры:\n\n- `data/` хранит содержимое таблицы;\n- `metadata/` хранит описание состояния таблицы;\n- JDBC catalog знает, где лежит актуальный metadata-файл.\n\nЭто и есть практическая разница между «таблица как логический объект» и «папка с parquet-файлами»." + }, + { + "cell_type": "markdown", + "id": "2b58a761", + "metadata": {}, + "source": "## 7. Почему `raw -> bronze` похоже на `stg -> ods`\n\nВ классическом DWH часто есть переход `stg -> ods`: сначала источник попадает в staging, потом появляется первая рабочая таблица, с которой удобно работать дальше.\n\nВ этом курсе аналог выглядит так:\n\n- `raw` — неизменяемый архив исходных файлов в object storage;\n- `bronze` — первая управляемая таблица на `Iceberg`, которую можно читать через каталог, осматривать через metadata-таблицы и позже запрашивать из `Trino`;\n- `silver` — следующий слой, где появятся нормализация, очистка и правила трансформаций.\n\nПоэтому `bronze` не равно `raw`:\n\n- raw нужен как воспроизводимая точка входа;\n- bronze нужен как управляемая таблица с каталогом и snapshot-ами;\n- bronze можно пересоздать из raw, если логика меняется." + }, + { + "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" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "e94508bf", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: прочитай taxi_zone_lookup.csv из raw-зоны через Spark\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "8e058aeb", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: создай таблицу lakehouse.bronze.taxi_zone_lookup\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "41c2c79a", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: загрузи данные в bronze lookup-таблицу\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "80b392e1", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: проверь результат через spark.table(...)\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "61e90959", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: посмотри snapshots для lakehouse.bronze.taxi_zone_lookup\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" + ] + }, + { + "cell_type": "markdown", + "id": "76d9061f", + "metadata": {}, + "source": [ + "## 10. Завершение\n", + "\n", + "Мы намеренно не удаляем `lakehouse.bronze.nyc_taxi_yellow`. Эта таблица понадобится в Модуле 5, где мы будем строить `silver`-слой.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "078b9876", + "metadata": {}, + "outputs": [], + "source": [ + "spark.stop()\n" + ] + } + ], + "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 +} \ No newline at end of file diff --git a/plans/README.md b/plans/README.md index 41e26b5..4efd315 100644 --- a/plans/README.md +++ b/plans/README.md @@ -16,3 +16,4 @@ | 1. Вход в стенд и базовая диагностика | [module-01-environment-and-basic-diagnostics.md](./module-01-environment-and-basic-diagnostics.md) | Ready for validation | `START_HERE.md`, `notebooks/01_environment_and_smoke_test.ipynb`, `src/spark/cluster_smoke.py` | 2026-03-06 | | 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 |