- Зачем: - необходимо дать студенту практический инструмент для понимания разделения 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/).
300 lines
13 KiB
Plaintext
300 lines
13 KiB
Plaintext
{
|
||
"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
|
||
} |