feat(notebooks): реализован ноутбук Модуля 2 «Ментальная модель Lakehouse»
- Зачем: - необходимо дать студенту практический инструмент для понимания разделения storage, catalog и compute. - Что: - создан notebooks/02_lakehouse_mental_model.ipynb с демонстрацией записи в Spark и инспекцией в PostgreSQL и MinIO. - в ноутбук добавлена таблица сравнения Lakehouse vs классические СУБД. - реализованы механизмы идемпотентности (CREATE OR REPLACE TABLE) и автоматической очистки ресурсов. - обновлен статус Модуля 2 в планах на "Ready for validation". - Проверка: - проверена структура JSON ноутбука, наличие всех демонстрационных и самостоятельных ячеек, а также корректность путей к MinIO (prefix warehouse/).
This commit is contained in:
@@ -0,0 +1,300 @@
|
||||
{
|
||||
"cells": [
|
||||
{
|
||||
"cell_type": "markdown",
|
||||
"metadata": {},
|
||||
"source": [
|
||||
"# Модуль 2. Ментальная модель Lakehouse: storage, catalog, compute\n",
|
||||
"\n",
|
||||
"В этом модуле мы на практике разберем, чем Lakehouse отличается от классических реляционных СУБД (PostgreSQL, Greenplum).\n",
|
||||
"В классической СУБД данные, метаданные и вычислитель жестко связаны внутри одной системы. В Lakehouse эти слои физически разделены:\n",
|
||||
"- **Storage (Хранилище)**: где лежат физические файлы (у нас это MinIO / S3).\n",
|
||||
"- **Catalog (Каталог)**: где хранится информация о таблицах и их расположении (у нас это PostgreSQL).\n",
|
||||
"- **Compute (Вычислитель)**: кто обрабатывает данные (у нас это Spark, а позже добавится Trino).\n",
|
||||
"\n",
|
||||
"| | PostgreSQL / Greenplum | Lakehouse |\n",
|
||||
"|---|---|---|\n",
|
||||
"| **Данные** | Внутри СУБД (pg_data) | Файлы в объектном хранилище (MinIO) |\n",
|
||||
"| **Метаданные** | Системные каталоги (pg_catalog) | Внешний каталог (PostgreSQL JDBC) |\n",
|
||||
"| **Вычислитель** | Тот же процесс СУБД | Независимые движки (Spark, Trino) |"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# Мы не указываем .master() явно, так как он и другие настройки Iceberg\n",
|
||||
"# подтягиваются из конфигурации spark-defaults.conf при старте сессии.\n",
|
||||
"from pyspark.sql import SparkSession\n",
|
||||
"\n",
|
||||
"spark = SparkSession.builder \\\n",
|
||||
" .appName(\"Module-02-Mental-Model\") \\\n",
|
||||
" .getOrCreate()\n",
|
||||
"\n",
|
||||
"# Spark выводит много INFO/WARN логов — убираем их, чтобы не засорять вывод ноутбука\n",
|
||||
"spark.sparkContext.setLogLevel(\"ERROR\")\n",
|
||||
"print(\"Spark сессия готова!\")"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "markdown",
|
||||
"metadata": {},
|
||||
"source": [
|
||||
"## 1. Запись таблицы (Compute)\n",
|
||||
"Мы используем Spark (вычислитель) для создания таблицы и записи данных."
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# Создаем namespace (логическую базу данных)\n",
|
||||
"spark.sql(\"CREATE NAMESPACE IF NOT EXISTS lakehouse.module_02\")\n",
|
||||
"\n",
|
||||
"# Создаем или перезаписываем таблицу и вставляем пару строк\n",
|
||||
"spark.sql(\"\"\"\n",
|
||||
"CREATE OR REPLACE TABLE lakehouse.module_02.demo_table (\n",
|
||||
" id int,\n",
|
||||
" name string,\n",
|
||||
" created_at timestamp\n",
|
||||
")\n",
|
||||
"USING iceberg\n",
|
||||
"\"\"\")\n",
|
||||
"\n",
|
||||
"spark.sql(\"\"\"\n",
|
||||
"INSERT INTO lakehouse.module_02.demo_table \n",
|
||||
"VALUES \n",
|
||||
" (1, 'Alice', current_timestamp()),\n",
|
||||
" (2, 'Bob', current_timestamp())\n",
|
||||
"\"\"\")\n",
|
||||
"\n",
|
||||
"# Прочитаем данные через Spark, чтобы убедиться, что они записались\n",
|
||||
"spark.sql(\"SELECT * FROM lakehouse.module_02.demo_table\").show()"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "markdown",
|
||||
"metadata": {},
|
||||
"source": [
|
||||
"## 2. Где живут метаданные? (Catalog)\n",
|
||||
"Таблица в Lakehouse — это не просто файлы. Чтобы Spark (или Trino) знал, где лежат данные и какая у них схема, используется Catalog. В нашем стенде это PostgreSQL.\n",
|
||||
"Посмотрим, какая запись появилась в базе данных каталога."
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"import psycopg2\n",
|
||||
"import pandas as pd\n",
|
||||
"import warnings\n",
|
||||
"\n",
|
||||
"# Подключаемся к PostgreSQL, который хранит метаданные Iceberg (наш Catalog)\n",
|
||||
"with psycopg2.connect(\n",
|
||||
" host=\"postgres-iceberg\",\n",
|
||||
" database=\"iceberg\",\n",
|
||||
" user=\"iceberg\",\n",
|
||||
" password=\"iceberg\"\n",
|
||||
") as conn:\n",
|
||||
" # Запрашиваем информацию о нашей таблице\n",
|
||||
" query = \"\"\"\n",
|
||||
" SELECT table_namespace, table_name, metadata_location, previous_metadata_location \n",
|
||||
" FROM iceberg_tables \n",
|
||||
" WHERE table_namespace = 'module_02'\n",
|
||||
" \"\"\"\n",
|
||||
"\n",
|
||||
" # Используем catch_warnings локально, чтобы скрыть UserWarning от pandas при чтении из DB-API\n",
|
||||
" with warnings.catch_warnings():\n",
|
||||
" warnings.simplefilter('ignore', UserWarning)\n",
|
||||
" df_catalog = pd.read_sql(query, conn)\n",
|
||||
"\n",
|
||||
"# Посмотрим, куда указывает catalog (обратите внимание на metadata_location)\n",
|
||||
"df_catalog"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "markdown",
|
||||
"metadata": {},
|
||||
"source": [
|
||||
"В поле `metadata_location` хранится путь к актуальному файлу метаданных таблицы в физическом хранилище.\n",
|
||||
"\n",
|
||||
"## 3. Где живут сами данные? (Storage)\n",
|
||||
"В качестве физического хранилища (Storage) мы используем MinIO — S3-совместимое объектное хранилище. Данные и файлы метаданных лежат там.\n",
|
||||
"Посмотрим на физические артефакты нашей таблицы."
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"import boto3\n",
|
||||
"\n",
|
||||
"# Подключаемся к MinIO\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",
|
||||
"bucket_name = 'lakehouse'\n",
|
||||
"# Путь (prefix) начинается с warehouse, так как это задано в настройках каталога\n",
|
||||
"prefix = 'warehouse/module_02/demo_table/'\n",
|
||||
"\n",
|
||||
"# Получаем список всех файлов (объектов), относящихся к нашей таблице\n",
|
||||
"response = s3.list_objects_v2(Bucket=bucket_name, Prefix=prefix)\n",
|
||||
"\n",
|
||||
"print(\"Файлы таблицы в MinIO:\\n\")\n",
|
||||
"contents = response.get('Contents', [])\n",
|
||||
"if not contents:\n",
|
||||
" print(\"Файлы не найдены. Проверьте prefix или убедитесь, что таблица создана.\")\n",
|
||||
"else:\n",
|
||||
" for obj in contents:\n",
|
||||
" key = obj['Key']\n",
|
||||
" size_kb = obj['Size'] / 1024\n",
|
||||
" if '/data/' in key:\n",
|
||||
" print(f\"[DATA] {key} ({size_kb:.2f} KB)\")\n",
|
||||
" elif '/metadata/' in key:\n",
|
||||
" print(f\"[META] {key} ({size_kb:.2f} KB)\")\n",
|
||||
" else:\n",
|
||||
" print(f\"[OTHER] {key} ({size_kb:.2f} KB)\")"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "markdown",
|
||||
"metadata": {},
|
||||
"source": [
|
||||
"**Что мы видим?**\n",
|
||||
"1. Папка `metadata/` содержит файлы `.json` (снимки состояния таблицы), `.avro` (манифесты, описывающие, какие файлы данных актуальны).\n",
|
||||
"2. Папка `data/` содержит сами данные в формате `.parquet`.\n",
|
||||
"\n",
|
||||
"> **Подумайте:** Что произойдет, если вы удалите файлы из MinIO, но не удалите запись в каталоге? А что если наоборот?"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "markdown",
|
||||
"metadata": {},
|
||||
"source": [
|
||||
"---\n",
|
||||
"\n",
|
||||
"## 4. Самостоятельное задание\n",
|
||||
"\n",
|
||||
"1. Создайте таблицу `lakehouse.module_02.my_first_table`.\n",
|
||||
"2. Запишите в неё несколько произвольных строк.\n",
|
||||
"3. Прочитайте таблицу через Spark.\n",
|
||||
"4. Выполните запросы к PostgreSQL и MinIO (аналогично примерам выше), чтобы найти `metadata_location` и физические `parquet` файлы вашей новой таблицы."
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# Ваш код: создание таблицы my_first_table и вставка данных\n",
|
||||
"\n"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# Ваш код: чтение данных через Spark\n",
|
||||
"\n"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# Ваш код: инспекция каталога (PostgreSQL)\n",
|
||||
"\n"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# Ваш код: инспекция хранилища (MinIO)\n",
|
||||
"\n"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "markdown",
|
||||
"metadata": {},
|
||||
"source": [
|
||||
"---\n",
|
||||
"\n",
|
||||
"## Checkpoint\n",
|
||||
"\n",
|
||||
"Убедитесь, что вы можете ответить на следующие вопросы:\n",
|
||||
"- Можете ли вы показать конкретные файлы данных и метаданных созданной таблицы в MinIO (через UI MinIO на порту 9001 или через код выше)?\n",
|
||||
"- Можете ли вы показать запись о таблице в PostgreSQL и объяснить, на что указывает поле `metadata_location`?\n",
|
||||
"- Что произойдёт, если удалить файлы из MinIO, но не трогать каталог (или наоборот)?\n",
|
||||
"- Понимаете ли вы теперь разницу между `storage` (MinIO) и `catalog` (PostgreSQL)?"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "markdown",
|
||||
"metadata": {},
|
||||
"source": [
|
||||
"---\n",
|
||||
"\n",
|
||||
"## 5. Очистка ресурсов (Cleanup)\n",
|
||||
"Выполните ячейку ниже, чтобы удалить созданные таблицы и namespace. Это необходимо, чтобы не оставлять \"мусор\" перед переходом к следующим модулям."
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# Удаляем таблицы (это удалит их и из каталога, и из хранилища)\n",
|
||||
"spark.sql(\"DROP TABLE IF EXISTS lakehouse.module_02.demo_table\")\n",
|
||||
"spark.sql(\"DROP TABLE IF EXISTS lakehouse.module_02.my_first_table\")\n",
|
||||
"\n",
|
||||
"# Удаляем namespace\n",
|
||||
"spark.sql(\"DROP NAMESPACE IF EXISTS lakehouse.module_02\")\n",
|
||||
"\n",
|
||||
"spark.stop()\n",
|
||||
"print(\"Очистка завершена, Spark сессия остановлена.\")"
|
||||
]
|
||||
}
|
||||
],
|
||||
"metadata": {
|
||||
"kernelspec": {
|
||||
"display_name": "Python 3",
|
||||
"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
|
||||
}
|
||||
+1
-1
@@ -14,4 +14,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) | Draft | `notebooks/02_lakehouse_mental_model.ipynb` | 2026-03-07 |
|
||||
| 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 |
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
# Модуль 2. Ментальная модель Lakehouse: storage, catalog, compute
|
||||
|
||||
**Статус:** `Draft`
|
||||
**Статус:** `Ready for validation`
|
||||
**Последнее обновление:** `2026-03-07`
|
||||
|
||||
## Цель
|
||||
|
||||
Reference in New Issue
Block a user