Files
mini-lakehouse-lab/notebooks/04_bronze_with_iceberg.ipynb
T
ddadmin 6567a44b74 feat(course): добавлен ноутбук модуля bronze с Iceberg
- Зачем:
  - нужен практический модуль 4, который переводит студента от raw-файлов к первой управляемой Iceberg-таблице.
- Что:
  - добавлен ноутбук notebooks/04_bronze_with_iceberg.ipynb с CTAS в bronze, разбором metadata tables и структурой хранения в MinIO.
  - добавлены пояснения про namespace, сравнение raw vs Iceberg и разбор snapshot-полей по итогам ревью.
  - обновлен plans/README.md строкой для модуля 4.
- Проверка:
  - синтаксис всех code-ячеек ноутбука проверен локальной компиляцией.
  - тесты пройдены пользователем.
2026-03-07 19:22:04 +03:00

692 lines
31 KiB
Plaintext
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.
{
"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
}