Files
mini-lakehouse-lab/notebooks/03_partitioning_and_schema_evolution.ipynb
T
2025-12-03 16:35:50 +03:00

164 lines
4.5 KiB
Plaintext
Executable File
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
{
"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
}