From 3763538aca6236e1e68d51297c94762d65c3519e Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Mar 2026 01:26:19 +0300 Subject: [PATCH] =?UTF-8?q?feat(notebooks):=20=D1=80=D0=B5=D0=B0=D0=BB?= =?UTF-8?q?=D0=B8=D0=B7=D0=BE=D0=B2=D0=B0=D0=BD=20=D0=BD=D0=BE=D1=83=D1=82?= =?UTF-8?q?=D0=B1=D1=83=D0=BA=20=D0=9C=D0=BE=D0=B4=D1=83=D0=BB=D1=8F=202?= =?UTF-8?q?=20=C2=AB=D0=9C=D0=B5=D0=BD=D1=82=D0=B0=D0=BB=D1=8C=D0=BD=D0=B0?= =?UTF-8?q?=D1=8F=20=D0=BC=D0=BE=D0=B4=D0=B5=D0=BB=D1=8C=20Lakehouse=C2=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - необходимо дать студенту практический инструмент для понимания разделения 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/). --- notebooks/02_lakehouse_mental_model.ipynb | 300 ++++++++++++++++++++++ plans/README.md | 2 +- plans/module-02-lakehouse-mental-model.md | 2 +- 3 files changed, 302 insertions(+), 2 deletions(-) create mode 100644 notebooks/02_lakehouse_mental_model.ipynb diff --git a/notebooks/02_lakehouse_mental_model.ipynb b/notebooks/02_lakehouse_mental_model.ipynb new file mode 100644 index 0000000..cd71428 --- /dev/null +++ b/notebooks/02_lakehouse_mental_model.ipynb @@ -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 +} \ No newline at end of file diff --git a/plans/README.md b/plans/README.md index 8d8aa87..a843ae2 100644 --- a/plans/README.md +++ b/plans/README.md @@ -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 | diff --git a/plans/module-02-lakehouse-mental-model.md b/plans/module-02-lakehouse-mental-model.md index ad62844..6aa9ceb 100644 --- a/plans/module-02-lakehouse-mental-model.md +++ b/plans/module-02-lakehouse-mental-model.md @@ -1,6 +1,6 @@ # Модуль 2. Ментальная модель Lakehouse: storage, catalog, compute -**Статус:** `Draft` +**Статус:** `Ready for validation` **Последнее обновление:** `2026-03-07` ## Цель