feat(module-07): реализован модуль по schema evolution и time travel
- Зачем: - обучение студентов безопасным изменениям таблиц и механизмам отката Iceberg. - Что: - создан notebooks/07_safe_table_changes.ipynb с демонстрацией ALTER TABLE, VERSION AS OF и rollback. - в ноутбук добавлена самостоятельная работа на демо-таблице и чекпоинт для самопроверки. - статус модуля 7 в планах обновлен до Ready for validation. - Проверка: - ручная верификация структуры ноутбука и синтаксиса Spark SQL/trino_query.
This commit is contained in:
@@ -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 <snapshot_id>` |\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 <snapshot_id>;`*\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_id>`? Что происходит с историей 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
|
||||
}
|
||||
Reference in New Issue
Block a user