- Зачем: - нужен практический модуль 4, который переводит студента от raw-файлов к первой управляемой Iceberg-таблице. - Что: - добавлен ноутбук notebooks/04_bronze_with_iceberg.ipynb с CTAS в bronze, разбором metadata tables и структурой хранения в MinIO. - добавлены пояснения про namespace, сравнение raw vs Iceberg и разбор snapshot-полей по итогам ревью. - обновлен plans/README.md строкой для модуля 4. - Проверка: - синтаксис всех code-ячеек ноутбука проверен локальной компиляцией. - тесты пройдены пользователем.
692 lines
31 KiB
Plaintext
692 lines
31 KiB
Plaintext
{
|
||
"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
|
||
} |