- Зачем: - обучение студентов безопасным изменениям таблиц и механизмам отката Iceberg. - Что: - создан notebooks/07_safe_table_changes.ipynb с демонстрацией ALTER TABLE, VERSION AS OF и rollback. - в ноутбук добавлена самостоятельная работа на демо-таблице и чекпоинт для самопроверки. - статус модуля 7 в планах обновлен до Ready for validation. - Проверка: - ручная верификация структуры ноутбука и синтаксиса Spark SQL/trino_query.
39 KiB
39 KiB
In [ ]:
import pandas as pd
from pyspark.sql import SparkSession
from trino.dbapi import connect
import trino
from IPython.display import display
spark = SparkSession.builder \
.appName("module-07-safe-table-changes") \
.getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
print("SparkSession создана.")
SILVER_TABLE = "lakehouse.silver.nyc_taxi_yellow"
DEMO_TABLE = "lakehouse.default.taxi_changes_demo"
In [ ]:
# Проверяем, что silver-таблица существует (Модуль 5 должен быть пройден)
assert spark.catalog.tableExists(SILVER_TABLE), f"Таблица {SILVER_TABLE} не найдена. Пройдите Модуль 5."
count = spark.table(SILVER_TABLE).count()
assert count > 0, f"Таблица {SILVER_TABLE} пуста."
print(f"Таблица {SILVER_TABLE} готова к работе. Строк: {count}")
In [ ]:
def trino_query(sql: str) -> pd.DataFrame:
"""Выполняет SQL-запрос в Trino и возвращает результат как Pandas DataFrame."""
conn = connect(
host="trino",
port=8080,
user="jupyter",
catalog="lakehouse"
)
cur = conn.cursor()
try:
cur.execute(sql)
rows = cur.fetchall()
columns = [desc[0] for desc in cur.description]
return pd.DataFrame(rows, columns=columns)
except trino.exceptions.ProgrammingError as e:
if "No nodes available to run query" in str(e):
print("Ошибка: Trino ещё не готов. Подождите пару минут после старта контейнеров.")
raise
# Запросы типа CREATE/DROP не возвращают строк
return pd.DataFrame([{"status": "Success"}])
finally:
conn.close()
print("Trino helper функция готова.")
In [ ]:
# Посмотрим snapshot-историю silver-таблицы
spark.sql(f"SELECT snapshot_id, committed_at, operation, summary FROM {SILVER_TABLE}.snapshots").show(truncate=False)
In [ ]:
# Дополнительная история с указанием предков (родительских snapshot-ов)
spark.sql(f"SELECT * FROM {SILVER_TABLE}.history").show(truncate=False)
In [ ]:
try:
spark.sql(f"ALTER TABLE {SILVER_TABLE} ADD COLUMNS (loaded_at TIMESTAMP)")
print("Колонка loaded_at добавлена.")
except Exception as e:
if "already exists" in str(e).lower():
print("Колонка loaded_at уже существует (повторный запуск ноутбука). Продолжаем.")
else:
raise
In [ ]:
# Проверим схему из Spark
spark.table(SILVER_TABLE).printSchema()
In [ ]:
# Посмотрим данные (должен быть NULL)
spark.sql(f"SELECT VendorID, loaded_at FROM {SILVER_TABLE} LIMIT 5").show()
In [ ]:
# Проверим схему из Trino
display(trino_query(f"DESCRIBE {SILVER_TABLE}"))
In [ ]:
# Прочитаем данные из Trino
display(trino_query(f"SELECT vendorid, loaded_at FROM {SILVER_TABLE} LIMIT 5"))
In [ ]:
# Создание демо-таблицы
spark.sql(f"""
CREATE OR REPLACE TABLE {DEMO_TABLE}
USING iceberg AS
SELECT VendorID, tpep_pickup_datetime, tpep_dropoff_datetime,
passenger_count, trip_distance, fare_amount, total_amount,
pickup_borough, pickup_zone
FROM {SILVER_TABLE}
LIMIT 1000
""")
count_1 = spark.table(DEMO_TABLE).count()
assert count_1 == 1000, f"Ожидалось 1000 строк, получено {count_1}"
print(f"Таблица создана. Строк: {count_1}")
In [ ]:
# Сохраним ID первого snapshot-а
df_snapshots = spark.sql(f"SELECT snapshot_id, committed_at FROM {DEMO_TABLE}.snapshots ORDER BY committed_at ASC").collect()
assert len(df_snapshots) == 1, f"Ожидался 1 snapshot, получено {len(df_snapshots)}"
snapshot_1 = df_snapshots[0]['snapshot_id']
committed_at_1 = df_snapshots[0]['committed_at']
print(f"Snapshot 1 ID: {snapshot_1} (создан: {committed_at_1})")
In [ ]:
# Добавим ещё 500 строк из Манхэттена (создаст второй snapshot)
spark.sql(f"""
INSERT INTO {DEMO_TABLE}
SELECT VendorID, tpep_pickup_datetime, tpep_dropoff_datetime,
passenger_count, trip_distance, fare_amount, total_amount,
pickup_borough, pickup_zone
FROM {SILVER_TABLE}
WHERE pickup_borough = 'Manhattan'
LIMIT 500
""")
count_2 = spark.table(DEMO_TABLE).count()
print(f"Добавлено 500 строк. Всего строк: {count_2}")
In [ ]:
# Сохраним ID второго snapshot-а
df_snapshots = spark.sql(f"SELECT snapshot_id, committed_at FROM {DEMO_TABLE}.snapshots ORDER BY committed_at ASC").collect()
snapshot_2 = df_snapshots[1]['snapshot_id']
committed_at_2 = df_snapshots[1]['committed_at']
print(f"Snapshot 2 ID: {snapshot_2} (создан: {committed_at_2})")
In [ ]:
# Чтение первого snapshot-а (начальная загрузка)
print(f"Читаем таблицу по состоянию Snapshot 1 ({snapshot_1}):")
spark.sql(f"""
SELECT count(*) AS row_count
FROM {DEMO_TABLE}
VERSION AS OF {snapshot_1}
""").show()
In [ ]:
# Чтение текущего состояния для сравнения
print("Читаем текущее состояние таблицы:")
spark.sql(f"""
SELECT count(*) AS row_count
FROM {DEMO_TABLE}
""").show()
In [ ]:
# Time travel через Trino к первому snapshot-у
print(f"Trino читает Snapshot 1 ({snapshot_1}):")
display(trino_query(f"""
SELECT count(*) AS row_count
FROM {DEMO_TABLE}
FOR VERSION AS OF {snapshot_1}
"""))
In [ ]:
# Time travel по времени (состояние на момент создания Snapshot 1)
spark.sql(f"""
SELECT count(*) AS row_count
FROM {DEMO_TABLE}
TIMESTAMP AS OF '{committed_at_1}'
""").show()
In [ ]:
# INSERT «плохих» данных (200 строк с отрицательным fare_amount)
spark.sql(f"""
INSERT INTO {DEMO_TABLE}
SELECT VendorID, tpep_pickup_datetime, tpep_dropoff_datetime,
passenger_count, trip_distance,
-99.99 AS fare_amount,
-99.99 AS total_amount,
'ERROR' AS pickup_borough,
'BAD_DATA' AS pickup_zone
FROM {SILVER_TABLE}
LIMIT 200
""")
# Убедимся, что таблица "испорчена"
spark.sql(f"""
SELECT count(*) as bad_rows
FROM {DEMO_TABLE}
WHERE fare_amount < 0
""").show()
count_3 = spark.table(DEMO_TABLE).count()
print(f"Всего строк сейчас: {count_3}")
In [ ]:
# Посмотрим на 3-й snapshot (нашу ошибку)
spark.sql(f"SELECT snapshot_id, operation FROM {DEMO_TABLE}.snapshots").show(truncate=False)
In [ ]:
# Выполняем rollback
spark.sql(f"""
CALL lakehouse.system.rollback_to_snapshot(
table => 'default.taxi_changes_demo',
snapshot_id => {snapshot_2}
)
""")
print(f"Rollback к snapshot {snapshot_2} завершен.")
In [ ]:
# Проверка восстановления
bad_rows = spark.sql(f"SELECT count(*) as bad_rows FROM {DEMO_TABLE} WHERE fare_amount < 0").collect()[0]['bad_rows']
current_count = spark.table(DEMO_TABLE).count()
print(f"Плохих строк (fare_amount < 0): {bad_rows}")
print(f"Всего строк: {current_count}")
In [ ]:
spark.sql(f"SELECT snapshot_id, operation FROM {DEMO_TABLE}.snapshots").show(truncate=False)
In [ ]:
# Переименуем колонку в демо-таблице
spark.sql(f"ALTER TABLE {DEMO_TABLE} RENAME COLUMN pickup_borough TO borough")
print("Колонка переименована в Spark.")
In [ ]:
# Посмотрим схему из Spark - колонка теперь borough
spark.table(DEMO_TABLE).printSchema()
In [ ]:
# Запрос с новым именем borough через Trino работает
display(trino_query(f"SELECT borough FROM {DEMO_TABLE} LIMIT 3"))
In [ ]:
# А запрос со старым именем pickup_borough упадет с ошибкой
try:
display(trino_query(f"SELECT pickup_borough FROM {DEMO_TABLE} LIMIT 3"))
except Exception as e:
print(f"Ошибка в Trino: {e}")
In [ ]:
# Пересоздаем демо-таблицу (это удалит всю ее snapshot-историю - Ошибка 1 в действии!)
spark.sql(f"""
CREATE OR REPLACE TABLE {DEMO_TABLE}
USING iceberg AS
SELECT VendorID, tpep_pickup_datetime, tpep_dropoff_datetime,
passenger_count, trip_distance, fare_amount, total_amount,
pickup_borough, pickup_zone
FROM {SILVER_TABLE}
LIMIT 1000
""")
print(f"Демо-таблица пересоздана начисто: {spark.table(DEMO_TABLE).count()} строк, 1 snapshot.")
In [ ]:
# Ваш код: ALTER TABLE ADD COLUMNS (quality_flag)
In [ ]:
# Ваш код: проверка из Spark и Trino
In [ ]:
# Ваш код: INSERT 300 строк из Brooklyn
In [ ]:
# Ваш код: snapshot-история
In [ ]:
# Ваш код: time travel — VERSION AS OF
In [ ]:
# Ваш код: INSERT 100 «плохих» строк
In [ ]:
# Ваш код: rollback_to_snapshot
In [ ]:
# Ваш код: проверка восстановления
In [ ]:
# Ваш код: snapshot-история silver
In [ ]:
spark.sql(f"DROP TABLE IF EXISTS {DEMO_TABLE}")
print("Демо-таблица удалена.")
In [ ]:
spark.stop()