diff --git a/notebooks/07_safe_table_changes.ipynb b/notebooks/07_safe_table_changes.ipynb new file mode 100644 index 0000000..196af74 --- /dev/null +++ b/notebooks/07_safe_table_changes.ipynb @@ -0,0 +1,904 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "# Модуль 7. Безопасная работа с таблицами: schema evolution и time travel\n", + "\n", + "**Цель:** Научиться безопасно изменять Iceberg-таблицы, использовать встроенные механизмы time travel (чтение предыдущих состояний) и rollback (восстановление после ошибок).\n", + "\n", + "В этом модуле мы сравним подходы классических DWH и Lakehouse к работе с изменениями:\n", + "\n", + "| | Классический DWH (Greenplum) | Lakehouse (Iceberg) |\n", + "|---|---|---|\n", + "| Добавить колонку | `ALTER TABLE ADD COLUMN` (мгновенно, NULL) | `ALTER TABLE ADD COLUMNS` (мгновенно, NULL) |\n", + "| Переименовать колонку | `ALTER TABLE RENAME COLUMN` | `ALTER TABLE RENAME COLUMN` |\n", + "| Откатить данные к вчерашнему состоянию | Восстановление из бэкапа (pg_dump / PITR) | `SELECT ... VERSION AS OF ` |\n", + "| Посмотреть историю изменений таблицы | Нет встроенного механизма | `SELECT * FROM table.snapshots` |\n", + "| Последствия ошибочной записи | Нужен бэкап или ручной откат | Rollback к предыдущему snapshot |\n", + "\n", + "**Ключевая идея:** Iceberg хранит полную историю изменений данных. Schema evolution (изменение структуры) и time travel (путешествие во времени) — встроенные инструменты формата, а не дополнительная инфраструктура.\n", + "\n", + "Как устроен этот ноутбук:\n", + "- Schema evolution демонстрируем на нашей рабочей таблице `silver` (это безопасная операция).\n", + "- Time travel и восстановление после ошибки — на отдельной демо-таблице (чтобы не рисковать данными для Модуля 8).\n", + "- Trino используется для проверки эффекта на downstream (системы, читающие данные после нас).\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-07-safe-table-changes\") \\\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_changes_demo\"\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Проверяем, что silver-таблица существует (Модуль 5 должен быть пройден)\n", + "assert spark.catalog.tableExists(SILVER_TABLE), f\"Таблица {SILVER_TABLE} не найдена. Пройдите Модуль 5.\"\n", + "count = spark.table(SILVER_TABLE).count()\n", + "assert count > 0, f\"Таблица {SILVER_TABLE} пуста.\"\n", + "print(f\"Таблица {SILVER_TABLE} готова к работе. Строк: {count}\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "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" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "*Если ячейки выше упали с `ImportError: trino` — нужно пересобрать образ Jupyter (в терминале: `docker compose build jupyter && docker compose up -d`).*\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Silver на месте, Trino готов. В этом модуле мы будем изменять таблицу и наблюдать последствия из обоих движков.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 2: Snapshot-ы — встроенная история изменений\n", + "\n", + "В Модуле 4 мы впервые увидели snapshot-ы bronze-таблицы. Теперь разберёмся глубже. Каждая операция с данными в Iceberg (INSERT, overwrite, delete) автоматически создаёт snapshot — «снимок» состояния таблицы. Snapshot фиксирует, какие data files составляли таблицу в этот момент. Это как коммит в git: можно вернуться к любому предыдущему состоянию.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Посмотрим snapshot-историю silver-таблицы\n", + "spark.sql(f\"SELECT snapshot_id, committed_at, operation, summary FROM {SILVER_TABLE}.snapshots\").show(truncate=False)\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Дополнительная история с указанием предков (родительских snapshot-ов)\n", + "spark.sql(f\"SELECT * FROM {SILVER_TABLE}.history\").show(truncate=False)\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Поля таблицы snapshots:\n", + "- `committed_at` — когда произошла операция.\n", + "- `operation` — тип операции (append, overwrite).\n", + "- `summary` — краткая статистика (added-data-files, total-records и т.д.).\n", + "\n", + "Если вы не перезапускали Модуль 5 через CREATE OR REPLACE, у silver будет 1 snapshot (от CTAS).\n", + "В Greenplum такой истории нет: после `INSERT` предыдущее состояние таблицы недоступно без бэкапа.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 3: Schema evolution — безопасное добавление колонки\n", + "\n", + "Добавляем колонку `loaded_at` к silver-таблице. В PostgreSQL `ALTER TABLE ADD COLUMN` — обычная операция. В Iceberg — аналогично, но с важным свойством: Iceberg хранит историю схем.\n", + "\n", + "Важно: `ALTER TABLE ADD COLUMNS` **НЕ создаёт** новый data snapshot. Это изменение метаданных (schema), а не данных: data files не перезаписываются, существующие строки логически получают NULL. Snapshot-ы фиксируют только операции с данными (INSERT, DELETE, overwrite).\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "try:\n", + " spark.sql(f\"ALTER TABLE {SILVER_TABLE} ADD COLUMNS (loaded_at TIMESTAMP)\")\n", + " print(\"Колонка loaded_at добавлена.\")\n", + "except Exception as e:\n", + " if \"already exists\" in str(e).lower():\n", + " print(\"Колонка loaded_at уже существует (повторный запуск ноутбука). Продолжаем.\")\n", + " else:\n", + " raise\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Проверим схему из Spark\n", + "spark.table(SILVER_TABLE).printSchema()\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Посмотрим данные (должен быть NULL)\n", + "spark.sql(f\"SELECT VendorID, loaded_at FROM {SILVER_TABLE} LIMIT 5\").show()\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Колонка добавлена. Данные не изменились — Iceberg не перезаписывал data files. \n", + "В рабочем сценарии `loaded_at` заполнялся бы при следующих INSERT-ах: `INSERT INTO ... SELECT ..., current_timestamp() AS loaded_at FROM ...`. Существующие строки остаются с NULL — это нормальная практика эволюции схемы. Заполнение существующих строк (UPDATE) — отдельная тема, требующая перезаписи файлов.\n", + "\n", + "Примечание: если вы повторно запустите Модуль 5 после Модуля 7, `CREATE OR REPLACE TABLE` пересоздаст silver без колонки `loaded_at`. Это ещё одна иллюстрация того, почему `CREATE OR REPLACE` опасен в рабочих сценариях — он стирает все изменения, включая schema evolution.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 4: Downstream-эффект — Trino видит изменение\n", + "\n", + "В Модуле 6 мы убедились: Spark и Trino читают одну таблицу через общий каталог. Schema evolution — ещё одно следствие этого дизайна: изменение схемы в Spark мгновенно видно в Trino. Не нужно «синхронизировать» или «обновлять» что-то на стороне Trino.\n", + "\n", + "*В DBeaver: `DESCRIBE lakehouse.silver.nyc_taxi_yellow;`*\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Проверим схему из Trino\n", + "display(trino_query(f\"DESCRIBE {SILVER_TABLE}\"))\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Прочитаем данные из Trino\n", + "display(trino_query(f\"SELECT vendorid, loaded_at FROM {SILVER_TABLE} LIMIT 5\"))\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Trino увидел новую колонку мгновенно. Spark изменил метаданные в PostgreSQL (JDBC catalog), Trino прочитал обновлённые метаданные. Decoupled compute + общий каталог в действии.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 5: Обзор операций schema evolution\n", + "\n", + "`ADD COLUMN` — самая безопасная операция. Но Iceberg поддерживает и другие:\n", + "\n", + "| Операция | Spark SQL | Безопасность | Влияние на downstream |\n", + "|---|---|---|---|\n", + "| Добавить колонку | `ALTER TABLE ADD COLUMNS (col type)` | Безопасно | Новая колонка с NULL, существующие запросы не ломаются |\n", + "| Переименовать колонку | `ALTER TABLE RENAME COLUMN old TO new` | Осторожно | Запросы с `SELECT old_name` перестают работать |\n", + "| Расширить тип | `ALTER TABLE ALTER COLUMN col TYPE bigint` | Безопасно | INT -> BIGINT: без потери данных |\n", + "| Сузить тип | — | Опасно | BIGINT -> INT: возможна потеря данных. Iceberg не поддерживает |\n", + "| Удалить колонку | `ALTER TABLE DROP COLUMN col` | Опасно | Запросы с `SELECT col` перестают работать |\n", + "\n", + "Далее в модуле мы продемонстрируем `RENAME COLUMN` на демо-таблице и покажем, как это ломает downstream-запросы.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 6: Подготовка демо-таблицы с несколькими snapshot-ами\n", + "\n", + "Для экспериментов с time travel и восстановлением (rollback) создаём отдельную демо-таблицу. Silver не трогаем — он нужен для финального Модуля 8. \n", + "\n", + "Мы создадим таблицу с небольшим подмножеством silver (1000 строк) и нарастим snapshot-историю через дополнительные `INSERT INTO`.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Создание демо-таблицы\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 1000\n", + "\"\"\")\n", + "\n", + "count_1 = spark.table(DEMO_TABLE).count()\n", + "assert count_1 == 1000, f\"Ожидалось 1000 строк, получено {count_1}\"\n", + "print(f\"Таблица создана. Строк: {count_1}\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Сохраним ID первого snapshot-а\n", + "df_snapshots = spark.sql(f\"SELECT snapshot_id, committed_at FROM {DEMO_TABLE}.snapshots ORDER BY committed_at ASC\").collect()\n", + "assert len(df_snapshots) == 1, f\"Ожидался 1 snapshot, получено {len(df_snapshots)}\"\n", + "snapshot_1 = df_snapshots[0]['snapshot_id']\n", + "committed_at_1 = df_snapshots[0]['committed_at']\n", + "\n", + "print(f\"Snapshot 1 ID: {snapshot_1} (создан: {committed_at_1})\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Добавим ещё 500 строк из Манхэттена (создаст второй snapshot)\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", + " WHERE pickup_borough = 'Manhattan'\n", + " LIMIT 500\n", + "\"\"\")\n", + "\n", + "count_2 = spark.table(DEMO_TABLE).count()\n", + "print(f\"Добавлено 500 строк. Всего строк: {count_2}\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Сохраним ID второго snapshot-а\n", + "df_snapshots = spark.sql(f\"SELECT snapshot_id, committed_at FROM {DEMO_TABLE}.snapshots ORDER BY committed_at ASC\").collect()\n", + "snapshot_2 = df_snapshots[1]['snapshot_id']\n", + "committed_at_2 = df_snapshots[1]['committed_at']\n", + "\n", + "print(f\"Snapshot 2 ID: {snapshot_2} (создан: {committed_at_2})\")\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Теперь у таблицы два snapshot-а: начальная загрузка (1000 строк) и добавление (500 строк). Это как два коммита в git. Каждый фиксирует конкретное состояние таблицы.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 7: Time travel — чтение предыдущих состояний\n", + "\n", + "Центральная возможность Iceberg: можно прочитать таблицу в том состоянии, в каком она была на момент любого snapshot-а. Это time travel. В Greenplum для этого пришлось бы восстанавливать бэкап целой базы или настраивать сложный PITR.\n", + "\n", + "#### 7a: Time travel через Spark\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Чтение первого snapshot-а (начальная загрузка)\n", + "print(f\"Читаем таблицу по состоянию Snapshot 1 ({snapshot_1}):\")\n", + "\n", + "spark.sql(f\"\"\"\n", + " SELECT count(*) AS row_count\n", + " FROM {DEMO_TABLE}\n", + " VERSION AS OF {snapshot_1}\n", + "\"\"\").show()\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Чтение текущего состояния для сравнения\n", + "print(\"Читаем текущее состояние таблицы:\")\n", + "\n", + "spark.sql(f\"\"\"\n", + " SELECT count(*) AS row_count\n", + " FROM {DEMO_TABLE}\n", + "\"\"\").show()\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Данные из первого snapshot-а не потеряны — оба состояния доступны одновременно. Iceberg хранит все data files; snapshot определяет, какие из них составляют таблицу в конкретный момент.\n", + "\n", + "#### 7b: Time travel через Trino\n", + "\n", + "Trino тоже умеет time travel. Тот же snapshot, та же таблица, тот же каталог.\n", + "Обратите внимание на синтаксис: в Spark — `VERSION AS OF`, в Trino — `FOR VERSION AS OF`.\n", + "\n", + "*В DBeaver: `SELECT count(*) FROM lakehouse.default.taxi_changes_demo FOR VERSION AS OF ;`*\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Time travel через Trino к первому snapshot-у\n", + "print(f\"Trino читает Snapshot 1 ({snapshot_1}):\")\n", + "\n", + "display(trino_query(f\"\"\"\n", + " SELECT count(*) AS row_count\n", + " FROM {DEMO_TABLE}\n", + " FOR VERSION AS OF {snapshot_1}\n", + "\"\"\"))\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "#### 7c: Time travel по времени\n", + "\n", + "Альтернативный способ — по timestamp. `TIMESTAMP AS OF` — Iceberg находит ближайший snapshot, существовавший на этот момент времени. Для точного отката лучше использовать `snapshot_id`.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Time travel по времени (состояние на момент создания Snapshot 1)\n", + "spark.sql(f\"\"\"\n", + " SELECT count(*) AS row_count\n", + " FROM {DEMO_TABLE}\n", + " TIMESTAMP AS OF '{committed_at_1}'\n", + "\"\"\").show()\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 8: Намеренная ошибка и восстановление (Rollback)\n", + "\n", + "Ключевой урок модуля. Мы намеренно внесём «плохие» данные, затем восстановим таблицу через rollback. Ошибки в Iceberg — не катастрофа, если понимаешь snapshot-историю.\n", + "\n", + "#### 8a: Вносим ошибку\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# INSERT «плохих» данных (200 строк с отрицательным fare_amount)\n", + "spark.sql(f\"\"\"\n", + " INSERT INTO {DEMO_TABLE}\n", + " SELECT VendorID, tpep_pickup_datetime, tpep_dropoff_datetime,\n", + " passenger_count, trip_distance,\n", + " -99.99 AS fare_amount,\n", + " -99.99 AS total_amount,\n", + " 'ERROR' AS pickup_borough,\n", + " 'BAD_DATA' AS pickup_zone\n", + " FROM {SILVER_TABLE}\n", + " LIMIT 200\n", + "\"\"\")\n", + "\n", + "# Убедимся, что таблица \"испорчена\"\n", + "spark.sql(f\"\"\"\n", + " SELECT count(*) as bad_rows \n", + " FROM {DEMO_TABLE} \n", + " WHERE fare_amount < 0\n", + "\"\"\").show()\n", + "\n", + "count_3 = spark.table(DEMO_TABLE).count()\n", + "print(f\"Всего строк сейчас: {count_3}\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Посмотрим на 3-й snapshot (нашу ошибку)\n", + "spark.sql(f\"SELECT snapshot_id, operation FROM {DEMO_TABLE}.snapshots\").show(truncate=False)\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "В Greenplum: 200 ошибочных строк в production-таблице -> звонок DBA, восстановление из бэкапа (если он есть и свежий), простой сервиса, или опасные ручные `DELETE`/`UPDATE` без возможности отката.\n", + "В Iceberg: rollback к предыдущему snapshot-у.\n", + "\n", + "#### 8b: Восстановление через rollback\n", + "Мы откатимся к snapshot 2, который мы сохранили в переменной `snapshot_2`.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Выполняем rollback\n", + "spark.sql(f\"\"\"\n", + " CALL lakehouse.system.rollback_to_snapshot(\n", + " table => 'default.taxi_changes_demo',\n", + " snapshot_id => {snapshot_2}\n", + " )\n", + "\"\"\")\n", + "print(f\"Rollback к snapshot {snapshot_2} завершен.\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Проверка восстановления\n", + "bad_rows = spark.sql(f\"SELECT count(*) as bad_rows FROM {DEMO_TABLE} WHERE fare_amount < 0\").collect()[0]['bad_rows']\n", + "current_count = spark.table(DEMO_TABLE).count()\n", + "\n", + "print(f\"Плохих строк (fare_amount < 0): {bad_rows}\")\n", + "print(f\"Всего строк: {current_count}\")\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Таблица восстановлена!\n", + "\n", + "#### 8c: Что произошло со snapshot-историей?\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "spark.sql(f\"SELECT snapshot_id, operation FROM {DEMO_TABLE}.snapshots\").show(truncate=False)\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Теперь в истории 4 записи. Четвёртый snapshot — это результат rollback, но его данные (набор файлов) в точности соответствуют snapshot 2.\n", + "\n", + "Rollback **не удаляет историю** и **не удаляет файлы**. Он создаёт новый snapshot, указывающий на данные предыдущего состояния. Аналогия:\n", + "\n", + "| | git | Iceberg |\n", + "|---|---|---|\n", + "| Откат с сохранением истории | `git revert` (новый коммит) | `rollback_to_snapshot` (новый snapshot) |\n", + "| Откат с удалением истории | `git reset --hard` | `CREATE OR REPLACE` (уничтожает snapshot-ы) |\n", + "\n", + "Старые data files (включая файлы с «плохими» строками) остаются в хранилище (MinIO) до процедуры `expire_snapshots` — это тема Модуля 8. Rollback не чистит storage, он только меняет указатель текущего snapshot-а.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 9: RENAME COLUMN и влияние на downstream\n", + "\n", + "Покажем, что `RENAME COLUMN` — безопасная операция для данных (они не перезаписываются), но опасная для downstream-запросов (они могут сломаться).\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Переименуем колонку в демо-таблице\n", + "spark.sql(f\"ALTER TABLE {DEMO_TABLE} RENAME COLUMN pickup_borough TO borough\")\n", + "print(\"Колонка переименована в Spark.\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Посмотрим схему из Spark - колонка теперь borough\n", + "spark.table(DEMO_TABLE).printSchema()\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Запрос с новым именем borough через Trino работает\n", + "display(trino_query(f\"SELECT borough FROM {DEMO_TABLE} LIMIT 3\"))\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# А запрос со старым именем pickup_borough упадет с ошибкой\n", + "try:\n", + " display(trino_query(f\"SELECT pickup_borough FROM {DEMO_TABLE} LIMIT 3\"))\n", + "except Exception as e:\n", + " print(f\"Ошибка в Trino: {e}\")\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Данные не потеряны — колонка переименована в метаданных. Но любой downstream-запрос, использовавший старое имя `pickup_borough`, сломается. В рабочем окружении перед `RENAME` нужно проверить, кто и как читает эту колонку.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 10: Типичные ошибки новичка\n", + "\n", + "**Ошибка 1: Слепой `CREATE OR REPLACE`.**\n", + "В Модулях 4 и 5 мы использовали `CREATE OR REPLACE TABLE ... AS SELECT` (CTAS) для создания bronze и silver. Это был осознанный компромисс: первичная загрузка, нет предыдущей истории, которую нужно сохранять. Но в рабочих сценариях `CREATE OR REPLACE` — **деструктивная операция**: она уничтожает ВСЮ snapshot-историю таблицы. Если у таблицы было 100 snapshot-ов, после `CREATE OR REPLACE` останется один. Time travel к предыдущим состояниям станет невозможен.\n", + "\n", + "**Ошибка 2: Изменение схемы без проверки downstream.**\n", + "Мы только что видели, как `RENAME COLUMN` ломает запросы со старым именем. `ALTER COLUMN TYPE` (расширение INT -> BIGINT) безопасен для данных, но downstream-системы могут интерпретировать новый тип по-другому и сломаться на своей стороне. Перед любым изменением схемы нужно понимать, кто читает эту таблицу.\n", + "\n", + "**Ошибка 3: Отсутствие проверки snapshot-истории перед деструктивной операцией.**\n", + "Перед любой сложной операцией, массово изменяющей данные (UPDATE/DELETE), полезно посмотреть текущее состояние: `SELECT * FROM table.snapshots`. Это занимает секунды и даёт понимание, какой snapshot станет точкой отката в случае проблемы. Привычка: посмотрел snapshot-ы -> выполнил операцию -> проверил результат -> убедился, что новый snapshot создан корректно.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 11: Самостоятельное задание\n", + "\n", + "Пересоздадим демо-таблицу с чистого листа. После демонстраций выше её схема изменена, а история содержит учебные эксперименты. Для самостоятельной работы нужна чистая таблица.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Пересоздаем демо-таблицу (это удалит всю ее snapshot-историю - Ошибка 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 1000\n", + "\"\"\")\n", + "print(f\"Демо-таблица пересоздана начисто: {spark.table(DEMO_TABLE).count()} строк, 1 snapshot.\")\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Выполни 5 задач на демо-таблице `lakehouse.default.taxi_changes_demo`.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "**Задача 1 (schema evolution).** Добавь колонку `quality_flag STRING` к демо-таблице через `ALTER TABLE ADD COLUMNS`. Проверь из Spark и из Trino, что колонка появилась и все значения NULL.\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "*DBeaver: `SELECT * FROM lakehouse.default.taxi_changes_demo LIMIT 10;`*\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: ALTER TABLE ADD COLUMNS (quality_flag)\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: проверка из Spark и Trino\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "**Задача 2 (создание snapshot).** Вставь (INSERT INTO) 300 строк из silver (с условием `pickup_borough = 'Brooklyn'`) в демо-таблицу. Проверь, что появился новый snapshot (всего их станет 2).\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: INSERT 300 строк из Brooklyn\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "**Задача 3 (time travel).** Посмотри snapshot-историю демо-таблицы. Прочитай таблицу в состоянии ДО добавления 300 строк (т.е. к snapshot 1) через `VERSION AS OF`. Сравни количество строк.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: snapshot-история\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: time travel — VERSION AS OF\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "**Задача 4 (recovery).** Сделай INSERT 100 строк с `fare_amount = -1` и `pickup_zone = 'STUDENT_ERROR'` в демо-таблицу. Убедись, что «плохие» строки появились. Затем выполни `rollback_to_snapshot` к snapshot-у ДО этого INSERT-а (т.е. к snapshot 2). Проверь, что таблица восстановлена.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: INSERT 100 «плохих» строк\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: rollback_to_snapshot\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: проверка восстановления\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "**Задача 5 (ответ).** Ответь текстом в ячейке ниже: чем `rollback_to_snapshot` отличается от `CREATE OR REPLACE TABLE ... AS SELECT * FROM table VERSION AS OF `? Что происходит с историей snapshot-ов в каждом случае?\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "Ваш ответ на Задачу 5: \n", + "...\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "*(Дополнительно, по желанию)*: Посмотри snapshot-историю silver-таблицы. Сколько snapshot-ов у неё? Какие операции их создали?\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: snapshot-история silver\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 12: Checkpoint\n", + "\n", + "Ответь для себя на вопросы:\n", + "1. Что такое snapshot в Iceberg и когда он создаётся?\n", + "2. Создаёт ли `ALTER TABLE ADD COLUMNS` новый data snapshot? Почему?\n", + "3. Чем `ADD COLUMN` отличается от `RENAME COLUMN` по влиянию на downstream?\n", + "4. Как прочитать предыдущее состояние таблицы? (назови два способа)\n", + "5. Что делает `rollback_to_snapshot`? Теряется ли при этом snapshot-история?\n", + "6. Почему `CREATE OR REPLACE` опаснее, чем `INSERT INTO` с последующим rollback?\n", + "7. Какой аналог time travel в классических СУБД (PostgreSQL/Greenplum) и почему он сложнее в использовании?\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Секция 13: Завершение\n", + "\n", + "Удаляем демо-таблицу — она больше не нужна. Silver остаётся для Модуля 8 (колонка `loaded_at` — демонстрационная; Модуль 8 не зависит от неё).\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "spark.sql(f\"DROP TABLE IF EXISTS {DEMO_TABLE}\")\n", + "print(\"Демо-таблица удалена.\")\n" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "**Итог модуля:**\n", + "- Schema evolution (`ADD COLUMN`) — безопасная metadata-only операция, которая мгновенно видна из всех движков через общий каталог.\n", + "- Snapshot-ы — автоматическая история изменений данных, встроенная в Iceberg.\n", + "- Time travel — чтение предыдущих состояний без бэкапов.\n", + "- Rollback — восстановление после ошибки с сохранением истории.\n", + "\n", + "В финальном Модуле 8 мы изучим базовое обслуживание таблиц: compaction (решение проблемы множества мелких файлов) и `expire_snapshots` (очистка старых файлов данных, оставшихся после time travel и rollback, которые мы научились применять в этом модуле).\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "spark.stop()\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 7376721..0a9ec15 100644 --- a/plans/README.md +++ b/plans/README.md @@ -19,5 +19,5 @@ | 4. Первая рабочая Iceberg-таблица и слой bronze | [module-04-bronze-with-iceberg.md](./module-04-bronze-with-iceberg.md) | Draft | `notebooks/04_bronze_with_iceberg.ipynb` | 2026-03-07 | | 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) | Draft | `notebooks/07_safe_table_changes.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 | diff --git a/plans/module-07-safe-table-changes.md b/plans/module-07-safe-table-changes.md index 79a8f60..77f2895 100644 --- a/plans/module-07-safe-table-changes.md +++ b/plans/module-07-safe-table-changes.md @@ -1,7 +1,7 @@ # Модуль 7. Безопасная работа с таблицами: schema evolution и time travel -**Статус:** `Draft` -**Последнее обновление:** `2026-03-07` +**Статус:** `Ready for validation` +**Последнее обновление:** `2026-03-08` ## Цель