diff --git a/notebooks/08_maintenance_and_final_lab.ipynb b/notebooks/08_maintenance_and_final_lab.ipynb new file mode 100644 index 0000000..9c45274 --- /dev/null +++ b/notebooks/08_maintenance_and_final_lab.ipynb @@ -0,0 +1,787 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "# Модуль 8. Базовое обслуживание таблиц и финальная практика\n", + "\n", + "**Цель:** Научиться базовому обслуживанию Iceberg-таблиц (compaction и expire_snapshots) и закрепить все навыки курса в финальной сквозной практике.\n", + "\n", + "**Prerequisites:** пройден Модуль 6, Модуль 7 рекомендуется.\n", + "\n", + "Две части модуля:\n", + "1. Базовое обслуживание: compaction (решение проблемы мелких файлов) и expire_snapshots (очистка истории).\n", + "2. Финальная практика: сквозной end-to-end сценарий.\n", + "\n", + "Зачем нужно обслуживание? Iceberg-таблица — не «чёрный ящик». После множества записей накапливаются мелкие файлы и старые snapshot-ы. Без обслуживания чтение замедляется (много мелких файлов), а хранилище растёт (старые данные не удаляются). Это знакомо из мира Greenplum: AppendOnly-таблицы тоже требуют VACUUM и REORGANIZE.\n", + "\n", + "| | Greenplum AppendOnly | Iceberg |\n", + "|---|---|---|\n", + "| Проблема | Мёртвые строки после UPDATE/DELETE | Мелкие файлы после множества INSERT |\n", + "| Компактификация | `ALTER TABLE ... REORGANIZE` | `rewrite_data_files` |\n", + "| Очистка устаревших версий | Нет встроенной истории версий | `expire_snapshots` |\n", + "| Очистка мёртвых строк/файлов | `VACUUM` | `expire_snapshots` (удаляет старые файлы) |\n", + "| Автоматизация в базовой конфигурации | Ручной запуск | Ручной запуск |\n", + "| В production | Расписание через cron/Airflow | Расписание через cron/Airflow |\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 1: Spark-сессия, Trino-клиент и проверка таблиц\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "import pandas as pd\n", + "from pyspark.sql import SparkSession\n", + "from trino.dbapi import connect\n", + "import trino\n", + "from IPython.display import display\n", + "\n", + "spark = SparkSession.builder \\\n", + " .appName(\"module-08-maintenance-and-lab\") \\\n", + " .getOrCreate()\n", + "\n", + "spark.sparkContext.setLogLevel(\"ERROR\")\n", + "print(\"SparkSession создана.\")\n", + "\n", + "SILVER_TABLE = \"lakehouse.silver.nyc_taxi_yellow\"\n", + "DEMO_TABLE = \"lakehouse.default.taxi_maintenance_demo\"\n", + "RAW_PATH = \"s3a://lakehouse/raw/nyc_taxi/yellow_tripdata_*.parquet\"\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "*Если ячейки выше упали с `ImportError: trino` — нужно пересобрать образ Jupyter (в терминале: `docker compose build jupyter && docker compose up -d`).*\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Проверяем, что silver-таблица существует (Модуль 5 должен быть пройден)\n", + "assert spark.catalog.tableExists(SILVER_TABLE), f\"Таблица {SILVER_TABLE} не найдена. Пройдите Модуль 5.\"\n", + "assert spark.table(SILVER_TABLE).count() > 0, f\"Таблица {SILVER_TABLE} пуста.\"\n", + "print(f\"Таблица {SILVER_TABLE} готова к работе.\")\n", + "\n", + "# Проверяем, что raw-данные доступны (Модуль 3)\n", + "try:\n", + " spark.read.parquet(RAW_PATH).limit(1).collect()\n", + " print(f\"Raw-данные доступны по пути: {RAW_PATH}\")\n", + "except Exception as e:\n", + " raise AssertionError(f\"Raw-данные не найдены по пути {RAW_PATH}. Вернись в Модуль 3 и выполни загрузку данных в MinIO.\") from e\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Helper для выполнения запросов к Trino\n", + "def trino_query(sql: str) -> pd.DataFrame:\n", + " \"\"\"Выполняет SQL-запрос в Trino и возвращает результат как Pandas DataFrame.\"\"\"\n", + " conn = connect(\n", + " host=\"trino\",\n", + " port=8080,\n", + " user=\"jupyter\",\n", + " catalog=\"lakehouse\"\n", + " )\n", + " cur = conn.cursor()\n", + " try:\n", + " cur.execute(sql)\n", + " rows = cur.fetchall()\n", + " columns = [desc[0] for desc in cur.description]\n", + " return pd.DataFrame(rows, columns=columns)\n", + " except trino.exceptions.ProgrammingError as e:\n", + " if \"No nodes available to run query\" in str(e):\n", + " print(\"Ошибка: Trino ещё не готов. Подождите пару минут после старта контейнеров.\")\n", + " raise\n", + " # Запросы типа CREATE/DROP не возвращают строк\n", + " return pd.DataFrame([{\"status\": \"Success\"}])\n", + " finally:\n", + " conn.close()\n", + "\n", + "print(\"Trino helper функция готова.\")\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Silver на месте, raw-данные доступны, Trino готов. В первой части модуля создадим демо-таблицу и научимся её обслуживать. Во второй — финальная практика.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 2: Проблема мелких файлов\n", + "\n", + "Каждый `INSERT INTO` в Iceberg создаёт новые data files. Если вставлять данные мелкими порциями (частые микробатчи, ручные INSERT-ы, инкрементальные загрузки), таблица накапливает множество мелких файлов. Это замедляет чтение: движку приходится открывать и читать каждый файл отдельно. \n", + "\n", + "Создадим демо-таблицу и сымитируем эту проблему: один изначальный `CTAS` (500 строк) и 8 маленьких `INSERT INTO` (по 100 строк).\n", + "*(Внимание: 8 INSERT-ов могут занять 1-2 минуты)*\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Создание демо-таблицы (создаст 1 файл)\n", + "spark.sql(f\"\"\"\n", + " CREATE OR REPLACE TABLE {DEMO_TABLE}\n", + " USING iceberg AS\n", + " SELECT VendorID, tpep_pickup_datetime, tpep_dropoff_datetime,\n", + " passenger_count, trip_distance, fare_amount, total_amount,\n", + " pickup_borough, pickup_zone\n", + " FROM {SILVER_TABLE}\n", + " LIMIT 500\n", + "\"\"\")\n", + "print(f\"Демо-таблица создана: {spark.table(DEMO_TABLE).count()} строк.\")\n", + "\n", + "# 8 мелких INSERT INTO (создадут 8 новых мелких файлов)\n", + "print(\"Запуск 8 INSERT-ов (может занять 1-2 минуты)...\")\n", + "for i in range(8):\n", + " spark.sql(f\"\"\"\n", + " INSERT INTO {DEMO_TABLE}\n", + " SELECT VendorID, tpep_pickup_datetime, tpep_dropoff_datetime,\n", + " passenger_count, trip_distance, fare_amount, total_amount,\n", + " pickup_borough, pickup_zone\n", + " FROM {SILVER_TABLE}\n", + " LIMIT 100\n", + " \"\"\")\n", + " print(f\"INSERT {i+1}/8 выполнен.\")\n", + "\n", + "print(f\"Итого строк: {spark.table(DEMO_TABLE).count()}\")\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Файловая статистика через metadata-таблицу files\n", + "files_df = spark.sql(f\"\"\"\n", + " SELECT file_path, file_format, record_count, file_size_in_bytes\n", + " FROM {DEMO_TABLE}.files\n", + "\"\"\")\n", + "print(f\"Количество data files: {files_df.count()}\")\n", + "files_df.show(truncate=40)\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Snapshot-история\n", + "spark.sql(f\"SELECT snapshot_id, committed_at, operation, summary FROM {DEMO_TABLE}.snapshots\").show(truncate=False)\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "**Результат:** 9 data files, 9 snapshot-ов. Каждый INSERT создал отдельный маленький файл. В production-сценарии после недель инкрементальных загрузок таблица может накопить сотни и тысячи мелких файлов.\n", + "\n", + "Для сравнения посмотрим на файловую статистику нашей основной `silver`-таблицы, которая была создана одним большим CTAS.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "silver_files = spark.sql(f\"SELECT file_path, record_count, file_size_in_bytes FROM {SILVER_TABLE}.files\")\n", + "print(f\"Silver: {silver_files.count()} data file(s)\")\n", + "silver_files.show(truncate=40)\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "`silver` имеет мало крупных файлов — проблемы мелких файлов нет. Compaction нужен не всегда, а после множества мелких записей. Далее мы исправим проблему на демо-таблице.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 3: Compaction — rewrite_data_files\n", + "\n", + "`rewrite_data_files` объединяет мелкие data files в меньшее количество крупных. \n", + "Аналогия: дефрагментация диска. Или `ALTER TABLE ... REORGANIZE` в Greenplum.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Сохраним статистику ДО compaction\n", + "files_before = spark.sql(f\"SELECT count(*) AS file_count, sum(file_size_in_bytes) AS total_bytes FROM {DEMO_TABLE}.files\").first()\n", + "print(f\"До compaction: {files_before['file_count']} файлов, {files_before['total_bytes']} байт\")\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Выполняем compaction\n", + "spark.sql(f\"CALL lakehouse.system.rewrite_data_files(table => 'default.taxi_maintenance_demo')\")\n", + "print(\"Compaction выполнен.\")\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Статистика ПОСЛЕ compaction\n", + "files_after = spark.sql(f\"SELECT count(*) AS file_count, sum(file_size_in_bytes) AS total_bytes FROM {DEMO_TABLE}.files\").first()\n", + "print(f\"После compaction: {files_after['file_count']} файлов, {files_after['total_bytes']} байт\")\n", + "print(f\"Было {files_before['file_count']} файлов → стало {files_after['file_count']} файлов\")\n", + "\n", + "# Проверим, что данные на месте\n", + "print(f\"Строк в таблице: {spark.table(DEMO_TABLE).count()}\")\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Данные не изменились — только реорганизованы файлы (чтение стало быстрее). \n", + "Compaction создал **новый snapshot** с новыми, более крупными файлами.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Посмотрим snapshot-историю: появился новый snapshot\n", + "spark.sql(f\"SELECT snapshot_id, committed_at, operation FROM {DEMO_TABLE}.snapshots\").show(truncate=False)\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Но старые мелкие файлы всё ещё лежат в MinIO: они нужны для time travel к старым snapshot-ам (например, чтобы мы могли запросить данные ДО compaction или ДО какого-то INSERT-а). Чтобы освободить место в хранилище, нужен `expire_snapshots`.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 4: expire_snapshots — очистка истории\n", + "\n", + "`expire_snapshots` удаляет старые snapshot-ы из метаданных. \n", + "**Важный компромисс:** после expire_snapshots time travel к удалённым snapshot-ам невозможен. Это необратимо. \n", + "В production нужно решить: сколько истории хранить? Обычно 10–30 дней. Для демо оставим только 2 последних snapshot-а.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Snapshot-история ДО expire\n", + "snapshots_before = spark.sql(f\"SELECT snapshot_id, committed_at, operation FROM {DEMO_TABLE}.snapshots\")\n", + "print(f\"Snapshot-ов до expire: {snapshots_before.count()}\")\n", + "snapshots_before.show(truncate=False)\n", + "\n", + "# Сохраняем ID самого старого snapshot для проверки time travel после expire\n", + "first_snapshot_id = snapshots_before.orderBy(\"committed_at\").first()[\"snapshot_id\"]\n", + "print(f\"Самый старый snapshot: {first_snapshot_id}\")\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Выполняем expire_snapshots (оставляем только 2 последних)\n", + "spark.sql(f\"\"\"\n", + " CALL lakehouse.system.expire_snapshots(\n", + " table => 'default.taxi_maintenance_demo',\n", + " retain_last => 2\n", + " )\n", + "\"\"\")\n", + "print(\"expire_snapshots выполнен.\")\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Snapshot-история ПОСЛЕ expire\n", + "snapshots_after = spark.sql(f\"SELECT snapshot_id, committed_at, operation FROM {DEMO_TABLE}.snapshots\")\n", + "print(f\"Snapshot-ов после expire: {snapshots_after.count()}\")\n", + "snapshots_after.show(truncate=False)\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Осталось всего 2 snapshot-а. Попробуем прочитать данные из удалённого (самого первого) snapshot-а:\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "try:\n", + " spark.sql(f\"\"\"\n", + " SELECT count(*) FROM {DEMO_TABLE}\n", + " VERSION AS OF {first_snapshot_id}\n", + " \"\"\").show()\n", + " print(\"Time travel сработал (snapshot не был удалён).\")\n", + "except Exception as e:\n", + " print(f\"Ожидаемая ошибка: {e}\")\n", + " print(\"Snapshot удалён — time travel невозможен.\")\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "**Вывод:** expire нужно запускать обдуманно. Слишком агрессивно (потеря страховки для rollback), слишком редко (рост метаданных и стоимости хранения).\n", + "\n", + "*Примечание: физическое удаление файлов из MinIO зависит от конфигурации. Иногда требуется дополнительно запускать процедуру `remove_orphan_files`, которая физически чистит файлы без привязки к активным snapshot-ам.*\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 5: Порядок обслуживания и типичные ошибки\n", + "\n", + "**Правильный порядок обслуживания Iceberg-таблицы:**\n", + "| Шаг | Операция | Что делает | Что будет, если пропустить |\n", + "|---|---|---|---|\n", + "| 1 | `rewrite_data_files` | Объединяет мелкие файлы в крупные | Чтение остаётся медленным |\n", + "| 2 | `expire_snapshots` | Удаляет старые snapshot-ы | Метаданные и хранилище растут |\n", + "\n", + "\n", + "Если сделать наоборот: expire удалит snapshot-ы, но мелкие файлы останутся (так как на них всё ещё ссылается текущий snapshot). Compaction затем создаст новые крупные файлы, а старые мелкие станут \"orphan\" (сиротами) и не удалятся автоматически.\n", + "\n", + "**Типичные ошибки:**\n", + "1. **Никогда не запускать compaction:** таблица обрастает тысячами мелких файлов, чтение деградирует в десятки раз.\n", + "2. **`expire_snapshots` сразу после ошибки (с `retain_last => 1`):** если записали плохие данные и запустили expire, откатиться назад (rollback) уже не выйдет. Правило: перед expire убедись, что текущее состояние корректно.\n", + "3. **Обслуживание на каждый INSERT:** compaction — тяжёлая операция. В production его делают по расписанию (например, 1 раз в сутки ночью).\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 6: Параллель с Greenplum AppendOnly\n", + "\n", + "Развернутое сравнение для студентов с опытом Greenplum:\n", + "\n", + "| Аспект | Greenplum AppendOnly | Iceberg |\n", + "|---|---|---|\n", + "| Как накапливается «мусор» | UPDATE/DELETE оставляют мёртвые строки в сегментах | Множество INSERT создают мелкие data files |\n", + "| Влияние на чтение | Сканирование мёртвых строк замедляет запросы | Открытие множества мелких файлов замедляет запросы |\n", + "| Компактификация | `ALTER TABLE ... REORGANIZE` | `rewrite_data_files` (объединение файлов) |\n", + "| Очистка | `VACUUM` (удаление мёртвых строк) | `expire_snapshots` (удаление старых snapshot-ов) |\n", + "| Автоматизация | `autovacuum` для heap-таблиц, ручной VACUUM для AO | Ручной запуск, автоматизация через Airflow/cron |\n", + "| Просмотр состояния | `pg_stat_all_tables`, `gp_toolkit.gp_bloat_diag` | `table.files`, `table.snapshots` |\n", + "| Time travel после очистки | Невозможен | Невозможен для удалённых snapshot-ов |\n", + "\n", + "Ключевое сходство: ни одна система не обслуживает себя сама. И там, и там это задача инженера. Ключевое отличие: Iceberg даёт встроенную историю (snapshot-ы) и контролируемую очистку (`retain_last`), чего нет в классическом DWH.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 7: Самостоятельное задание — обслуживание\n", + "\n", + "**Задача 1 (анализ silver).** Посмотри файловую статистику `silver`-таблицы (`lakehouse.silver.nyc_taxi_yellow`): количество data files, средний размер файла, количество snapshot-ов. Сравни с демо-таблицей до compaction. Ответь: нужна ли silver compaction? При каких условиях потребовалась бы?\n", + "\n", + "**Задача 2 (практика).** \n", + "- Создай таблицу `lakehouse.default.taxi_compaction_practice` через CTAS (500 строк из `silver`). \n", + "- Выполни 5 INSERT INTO по 200 строк. \n", + "- Посмотри файловую статистику ДО.\n", + "- Выполни compaction и expire_snapshots (retain_last => 1). \n", + "- Проверь результат ПОСЛЕ. \n", + "- Удали таблицу (`DROP TABLE`).\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: файловая статистика silver (файлы, размеры, snapshot-ы)\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Ваш ответ на Задачу 1: нужна ли compaction для silver? При каких условиях потребовалась бы? \n", + "...\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: создание taxi_compaction_practice + 5 INSERT\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: файловая статистика до compaction\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: compaction + expire_snapshots\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: файловая статистика после + DROP TABLE\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 8: Удаление демо-таблицы\n", + "\n", + "Демо-таблица для обслуживания больше не нужна.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "spark.sql(f\"DROP TABLE IF EXISTS {DEMO_TABLE}\")\n", + "print(\"Демо-таблица удалена.\")\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 9: Финальная практика — сквозной end-to-end сценарий\n", + "\n", + "Финальная практика курса. Самостоятельно пройди полный цикл Lakehouse-пайплайна: от raw-данных до обслуживания таблицы. Все навыки из Модулей 1–8 в одном сценарии.\n", + "\n", + "Работаем в отдельном namespace `lakehouse.final`, чтобы не затрагивать основные таблицы.\n", + "\n", + "**Сценарий (8 шагов):**\n", + "\n", + "**Шаг 1 (namespace).** Создай namespace `lakehouse.final`.\n", + "\n", + "**Шаг 2 (bronze).** Прочитай raw parquet из MinIO (`s3a://lakehouse/raw/nyc_taxi/yellow_tripdata_*.parquet`, LIMIT 5000 строк) и создай таблицу `lakehouse.final.trips_bronze` через CTAS. Это bronze: данные «как есть».\n", + "\n", + "**Шаг 3 (silver).** Создай `lakehouse.final.trips_silver` из bronze с тремя трансформациями:\n", + "- Фильтр: `fare_amount >= 0`\n", + "- Нормализация: `passenger_count` → INT с обработкой NULL (`coalesce` + `cast`)\n", + "- Обогащение: вычисляемая колонка `trip_duration_minutes` (разница между dropoff и pickup в минутах)\n", + "\n", + "**Шаг 4 (Trino).** Прочитай `lakehouse.final.trips_silver` из Trino. Сравни count с результатом из Spark.\n", + "\n", + "**Шаг 5 (schema evolution).** Добавь колонку `processed_at TIMESTAMP` к `trips_silver` через `ALTER TABLE ADD COLUMNS`. Проверь из Trino (`DESCRIBE`), что колонка видна.\n", + "\n", + "**Шаг 6 (snapshot-ы).** Выполни 3 `INSERT INTO trips_silver` (например, по 300 строк из `trips_bronze`). Каждый INSERT должен включать трансформации из шага 3, плюс `current_timestamp() AS processed_at`. *Подсказка: если не указать `processed_at`, Spark выдаст ошибку несовпадения числа колонок; альтернативный вариант — использовать явный список target-колонок в `INSERT INTO trips_silver (col1, col2, ...)`.* Посмотри файловую статистику и snapshot-ы.\n", + "\n", + "**Шаг 7 (обслуживание).** Выполни compaction (`rewrite_data_files`), затем `expire_snapshots` (retain_last => 1). Проверь, что файлов стало меньше, а snapshot-ов остался 1.\n", + "\n", + "**Шаг 8 (cleanup).** Удали `trips_silver`, `trips_bronze` и namespace `lakehouse.final`.\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 1: CREATE NAMESPACE\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 2: прочитай raw и создай trips_bronze (LIMIT 5000)\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 2: проверка trips_bronze (count, printSchema)\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 3: создай trips_silver с трансформациями\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 3: проверка trips_silver (count, проверка фильтра)\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "*DBeaver: `SELECT count(*) FROM lakehouse.final.trips_silver;`*\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 4: чтение trips_silver из Trino\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 5: ALTER TABLE ADD COLUMNS\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "*DBeaver: `DESCRIBE lakehouse.final.trips_silver;`*\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 5: проверка из Trino (DESCRIBE)\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 6: 3 INSERT INTO trips_silver\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 6: файловая статистика и snapshot-ы\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 7: compaction + expire_snapshots\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 7: проверка результата\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Шаг 8: cleanup (DROP TABLE, DROP NAMESPACE)\n", + "\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 10: Финальный checkpoint\n", + "\n", + "**Часть A — Обслуживание таблиц (Модуль 8):**\n", + "1. Что такое проблема мелких файлов и когда она возникает?\n", + "2. Что делает `rewrite_data_files`? Удаляет ли он старые файлы?\n", + "3. Что делает `expire_snapshots`? Какой компромисс он создаёт?\n", + "4. В каком порядке нужно выполнять compaction и expire_snapshots? Почему?\n", + "5. Как обслуживание Iceberg-таблиц соотносится с VACUUM/REORGANIZE в Greenplum?\n", + "\n", + "**Часть B — Весь курс (итоговые вопросы):**\n", + "1. Назови три роли в архитектуре Lakehouse и какие компоненты стенда их выполняют.\n", + "2. Почему Spark может записать данные, а Trino прочитать ту же таблицу без копирования?\n", + "3. Чем raw отличается от bronze? Чем bronze отличается от silver?\n", + "4. Что такое snapshot в Iceberg? Чем он отличается от бэкапа в PostgreSQL?\n", + "5. Назови три безопасные и три опасные операции с Iceberg-таблицей.\n", + "6. Зачем нужны compaction и expire_snapshots в production-сценарии?\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 11: Завершение\n", + "\n", + "**Итог модуля:**\n", + "- Compaction (`rewrite_data_files`) — решает проблему мелких файлов, не меняя данные (чтение становится быстрее).\n", + "- `expire_snapshots` — очищает старые snapshot-ы (сокращая размер метаданных и хранилища), но лишает возможности time travel.\n", + "- Порядок: compaction → expire_snapshots.\n", + "- Обслуживание Iceberg похоже на обслуживание Greenplum AppendOnly: обе системы требуют периодической заботы.\n", + "\n", + "**Итог курса. Что ты теперь умеешь:**\n", + "1. Поднимать и диагностировать локальный Lakehouse-стенд (Модуль 1).\n", + "2. Объяснять роли storage, catalog, compute (Модуль 2).\n", + "3. Загружать raw-данные и проверять схему (Модуль 3).\n", + "4. Создавать и заполнять Iceberg-таблицы (Модуль 4).\n", + "5. Строить воспроизводимый поток raw → bronze → silver (Модуль 5).\n", + "6. Читать одну таблицу из Spark и Trino (Модуль 6).\n", + "7. Безопасно изменять таблицы: schema evolution, time travel, rollback (Модуль 7).\n", + "8. Обслуживать таблицы: compaction, expire_snapshots (Модуль 8).\n", + "\n", + "Это базовый набор навыков для безопасной работы в Lakehouse-стенде. Дальше — production: партиционирование, MERGE, streaming, Airflow, governance. Но фундамент заложен.\n", + "\n", + "**Курс завершен! Поздравляем! 🎉**\n", + "\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "spark.stop()\n", + "\n" + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": "Python 3", + "language": "python", + "name": "python3" + }, + "language_info": { + "name": "python" + } + }, + "nbformat": 4, + "nbformat_minor": 4 +} \ No newline at end of file diff --git a/plans/README.md b/plans/README.md index 0a9ec15..f2c9eed 100644 --- a/plans/README.md +++ b/plans/README.md @@ -20,4 +20,4 @@ | 5. Слой silver и воспроизводимые трансформации | [module-05-silver-layer.md](./module-05-silver-layer.md) | Ready for validation | `notebooks/05_silver_layer.ipynb` | 2026-03-07 | | 6. Одна таблица, два движка: Spark и Trino | [module-06-spark-and-trino-on-same-table.md](./module-06-spark-and-trino-on-same-table.md) | Ready for validation | `jupyter/Dockerfile`, `START_HERE.md`, `docs/stack_reference.md`, `notebooks/06_spark_and_trino_on_same_table.ipynb` | 2026-03-07 | | 7. Безопасная работа с таблицами: schema evolution и time travel | [module-07-safe-table-changes.md](./module-07-safe-table-changes.md) | Ready for validation | `notebooks/07_safe_table_changes.ipynb` | 2026-03-08 | -| 8. Базовое обслуживание таблиц и финальная практика | [module-08-maintenance-and-final-lab.md](./module-08-maintenance-and-final-lab.md) | Draft | `notebooks/08_maintenance_and_final_lab.ipynb` | 2026-03-07 | +| 8. Базовое обслуживание таблиц и финальная практика | [module-08-maintenance-and-final-lab.md](./module-08-maintenance-and-final-lab.md) | Ready for validation | `notebooks/08_maintenance_and_final_lab.ipynb` | 2026-03-08 | diff --git a/plans/module-08-maintenance-and-final-lab.md b/plans/module-08-maintenance-and-final-lab.md index b37e823..8e25d42 100644 --- a/plans/module-08-maintenance-and-final-lab.md +++ b/plans/module-08-maintenance-and-final-lab.md @@ -1,7 +1,7 @@ # Модуль 8. Базовое обслуживание таблиц и финальная практика -**Статус:** `Draft` -**Последнее обновление:** `2026-03-07` +**Статус:** `Ready for validation` +**Последнее обновление:** `2026-03-08` ## Цель