{ "cells": [ { "cell_type": "markdown", "metadata": {}, "source": [ "# Лаба 3: партиционирование и эволюция схемы", "", "В этом ноутбуке мы посмотрим, как Iceberg работает с партиционированными таблицами и эволюцией схемы поверх общего каталога `lakehouse`.", "", "Перед началом убедись, что стенд запущен (`docker compose up -d`) и Spark-кластер доступен." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "from pyspark.sql import SparkSession", "", "spark = (", " SparkSession.builder", " .appName(\"lab3-partitioning-schema-evolution\")", " .master(\"spark://spark-master:7077\")", " .getOrCreate()", ")" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## Часть 1. Партиционированная таблица", "", "Для начала создадим партиционированную Iceberg-таблицу `lakehouse.default.partition_demo` с помощью готового SQL-скрипта из `src/spark/partitioned_table_demo.sql`." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "from pathlib import Path", "", "sql_path = Path(\"/opt/src/spark/partitioned_table_demo.sql\")", "sql_text = sql_path.read_text(encoding=\"utf-8\")", "", "for statement in sql_text.split(\";\"):", " stmt = statement.strip()", " if stmt:", " spark.sql(stmt)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "Посмотрим, какие данные записаны по датам, и обсудим партиционирование по `event_date`." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "spark.sql(\"\"\"", "SELECT", " event_date,", " COUNT(*) AS cnt,", " SUM(amount) AS total_amount", "FROM lakehouse.default.partition_demo", "GROUP BY event_date", "ORDER BY event_date", "\"\"\").show()" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "spark.sql(\"\"\"", "SELECT *", "FROM lakehouse.default.partition_demo", "WHERE event_date = DATE '2024-01-01'", "ORDER BY user_id", "\"\"\").show()" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## Часть 2. Эволюция схемы", "", "Теперь посмотрим на эволюцию схемы: создадим таблицу, добавим колонку и вставим новые строки с дополнительными данными.", "", "Для подготовки таблицы используем скрипт `src/spark/schema_evolution_demo.sql`." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "sql_path = Path(\"/opt/src/spark/schema_evolution_demo.sql\")", "sql_text = sql_path.read_text(encoding=\"utf-8\")", "", "for statement in sql_text.split(\";\"):", " stmt = statement.strip()", " if stmt:", " spark.sql(stmt)" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "spark.sql(\"DESCRIBE TABLE lakehouse.default.schema_evolution_demo\").show(truncate=False)" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "spark.sql(\"\"\"", "SELECT *", "FROM lakehouse.default.schema_evolution_demo", "ORDER BY id", "\"\"\").show(truncate=False)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "Обрати внимание, что старые строки имеют `NULL` в колонке `metadata`, а новые — заполненное значение.", "Iceberg хранит историю снапшотов и позволяет эволюцию схемы без сложных миграций." ] } ], "metadata": { "kernelspec": { "display_name": "Python 3", "language": "python", "name": "python3" }, "language_info": { "name": "python" } }, "nbformat": 4, "nbformat_minor": 5 }