commit 41bbdecba577c63fe92bf817acf752708883abea Author: Dmitrii Date: Wed Dec 3 16:35:50 2025 +0300 первый commit diff --git a/.gitignore b/.gitignore new file mode 100755 index 0000000..b3ea1eb --- /dev/null +++ b/.gitignore @@ -0,0 +1,23 @@ +.ipynb_checkpoints +__pycache__/ + +# Python bytecode / build artifacts +*.py[cod] +*.pyo +*.pyd +*.egg-info/ +*.egg + +# Virtual environments +.venv/ +venv/ + +# Editor / IDE settings +.idea/ +.vscode/ + +# Local config +.env + +# Logs +*.log diff --git a/AGENTS.md b/AGENTS.md new file mode 100755 index 0000000..7b47720 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,56 @@ +# Repository Guidelines + +This repository contains a small teaching Lakehouse stack (Spark + Trino + Iceberg + MinIO) intended for demos and mentoring. + +## Project Structure & Module Organization + +- `docker-compose.yml`: orchestrates Spark, Trino, MinIO, PostgreSQL, and Jupyter. +- `spark/`: Spark image (`Dockerfile`) and `spark-defaults.conf`. +- `jupyter/`: Jupyter image (`Dockerfile`) based on the Spark image. +- `trino/`: catalog config, e.g. `trino/catalog/lakehouse.properties`. +- `src/`: PySpark and SQL examples split by engine (`src/spark/...`, `src/trino/...`). +- `notebooks/`: demo notebooks (mounted into Jupyter at `/opt/work`). + +Keep new examples in `src/` or `notebooks/`, and avoid mixing configuration and code. + +## Build, Test, and Development Commands + +Run from the repo root: + +- `docker compose build`: build custom Spark and Jupyter images. +- `docker compose up -d`: start the full stack in the background. +- `docker compose ps`: check container status. +- `docker compose down -v`: stop the stack and remove volumes (for a clean slate). + +Use `docker compose logs -f ` when debugging (`spark-master`, `trino`, `minio`, etc.). + +## Coding Style & Naming Conventions + +- Python: PEP 8, 4-space indentation, snake_case for functions, lower_snake_case for files (e.g. `spark_join_demo.py`). +- SQL: uppercase keywords, `schema.table` naming, short English identifiers; comments may be in Russian. +- Compose/Docker: service names kebab-case (`spark-master`), env vars UPPER_SNAKE_CASE. +- Prefer small, focused examples; reuse helpers from `src/` in notebooks where possible. + +## Testing Guidelines + +There is no formal automated test suite yet. Validate changes by: + +- building and starting the stack, then +- running `src/spark/cluster_smoke.py` or the SQL in `src/spark/iceberg_demo.sql`, +- opening `notebooks/spark-basic-test.ipynb` in Jupyter and checking it runs end-to-end. + +If you add tests, prefer `pytest` under `tests/` and mark slow, integration-heavy tests clearly. + +## Commit & Pull Request Guidelines + +- Commits: short, descriptive subject in Russian or English, present tense (e.g. `Спарк запускается, создаётся тестовая таблица`, `Add Iceberg join demo`). +- Keep changes small and focused; update `README.md` when behavior, ports, or images change. +- PRs should describe the problem, the solution, and any impact on local setup; include screenshots (Trino UI, MinIO, Jupyter) when UI changes are relevant. + +## Agent-Specific Instructions + +When editing as an automated agent: + +- Respect this file and avoid large refactors without an explicit request. +- Preserve Russian user-facing text unless the change is explicitly about translation. +- Prefer minimal changes that keep the demo simple and robust for newcomers. diff --git a/HOWTO.md b/HOWTO.md new file mode 100755 index 0000000..4be75ef --- /dev/null +++ b/HOWTO.md @@ -0,0 +1,298 @@ +# HOWTO: Учебные лабораторные по Lakehouse-стенду + +Этот файл описывает пошаговые учебные сценарии (лабораторные работы), которые можно выполнять поверх стенда из `docker-compose.yml`. + +Перед началом работ смотри разделы «Сборка и запуск» и «Доступ к сервисам» в `README.md` — там описано, как поднять стенд и на каких портах доступны Trino, MinIO, Spark и Jupyter. + +Лабораторки опираются на примеры в `src/`: + +- `src/spark/cluster_smoke.py` — проверка, что Spark‑кластер жив. +- `src/spark/iceberg_demo.sql` — первая Iceberg‑таблица в Spark. +- `src/spark/iceberg_smoke.py` — Spark создаёт Iceberg‑таблицу и читает её. +- `src/trino/iceberg_smoke.sql` — Trino читает таблицу, созданную в Spark. + +Все команды ниже выполняются из корня репозитория. + +Краткая карта лабораторных: + +- **Лаба 0** — стенд поднят, Spark‑кластер жив. +- **Лаба 1** — первая Iceberg‑таблица в Spark. +- **Лаба 2** — общий каталог Spark ↔ Trino. +- **Лаба 3** — партиционирование и эволюция схемы. +- **Лаба 4** — мини‑ETL поверх Lakehouse (эскиз). + +--- + +## Лаба 0. Стенд поднят, Spark‑кластер жив + +**Цель** + +- Убедиться, что все контейнеры поднялись. +- Проверить, что Spark‑кластер (master + workers) работает. + +**Предусловия** + +- Установлены Docker и Docker Compose. +- Репозиторий склонирован локально. + +**Шаги** + +1. Собрать и поднять стенд: + + ```bash + docker compose build + docker compose up -d + ``` + +2. Проверить статус контейнеров: + + ```bash + docker compose ps + ``` + + Ожидаем, что `spark-master`, `spark-worker-1`, `spark-worker-2` в статусе `Up`. + +3. Запустить smoke‑скрипт кластера из контейнера `spark-master`: + + ```bash + docker compose exec spark-master \ + /opt/spark/bin/spark-submit /opt/src/spark/cluster_smoke.py + ``` + +4. Открыть Web UI Spark: + +- `http://localhost:8080` — мастер. +- `http://localhost:8081` и `http://localhost:8082` — воркеры. + +**Ожидаемый результат** + +- `cluster_smoke.py` выполняется без ошибок. +- В Web UI видно приложение, прошедшее через кластер. + +--- + +## Лаба 1. Первая Iceberg‑таблица из Spark + +**Цель** + +- Создать Iceberg‑таблицу с помощью Spark SQL. +- Посмотреть файлы таблицы в MinIO (`warehouse/default/demo_tbl/...`). + +**Предусловия** + +- Стенд запущен (лаба 0 выполнена). + +**Шаги** + +1. Подключиться к Spark SQL в контейнере: + + ```bash + docker compose exec -it spark-master \ + /opt/spark/bin/spark-sql + ``` + +2. Выполнить учебный SQL‑скрипт: + +- Внутри интерактивной сессии `spark-sql`: + + ```sql + :r /opt/src/spark/iceberg_demo.sql + ``` + +- Либо одним вызовом (без интерактивного режима): + + ```bash + docker compose exec spark-master \ + /opt/spark/bin/spark-sql -f /opt/src/spark/iceberg_demo.sql + ``` + +3. Зайти в MinIO Console: + +- Адрес: `http://localhost:9001` +- Учётные данные по умолчанию: `minioadmin / minioadmin` (см. `docker-compose.yml`). + +4. Найти файлы таблицы: + +- Бакет `lakehouse`. +- Префикс `warehouse/default/demo_tbl/`. +- Обратить внимание на структуру Iceberg: каталоги `metadata/`, `data/` и т.д. + +**Ожидаемый результат** + +- В Spark запрос + + ```sql + SELECT * FROM lakehouse.default.demo_tbl; + ``` + + возвращает данные. + +- В MinIO видна структура Iceberg‑таблицы: служебные файлы и файлы данных. + +--- + +## Лаба 2. Общий каталог Spark ↔ Trino + +**Цель** + +- Показать, что Spark и Trino используют общий Iceberg‑каталог (метаданные в Postgres, данные в MinIO). +- Создать таблицу из Spark и прочитать её через Trino. + +**Предусловия** + +- Стенд запущен. +- Лаба 1 не обязательна, но полезна для понимания структуры файлов. + +**Шаги** + +1. Создать таблицу и записать данные из Spark (PySpark‑скрипт): + + ```bash + docker compose exec spark-master \ + /opt/spark/bin/spark-submit /opt/src/spark/iceberg_smoke.py + ``` + +2. Убедиться в Spark, что таблица существует: + + ```bash + docker compose exec -it spark-master \ + /opt/spark/bin/spark-sql + ``` + + В интерактивной сессии: + + ```sql + USE lakehouse.default; + SHOW TABLES; + SELECT * FROM spark_trino_smoke; + ``` + +3. Прочитать ту же таблицу из Trino: + +- Вариант через заранее скопированный SQL‑файл (как в README): + + ```bash + docker compose cp src/trino/iceberg_smoke.sql trino:/tmp/ + docker compose exec trino trino --file /tmp/iceberg_smoke.sql + ``` + +- Либо интерактивно внутри Trino CLI: + + ```bash + docker compose exec -it trino trino --catalog lakehouse + ``` + + Внутри CLI: + + ```sql + USE lakehouse.default; + SHOW TABLES; + SELECT * FROM spark_trino_smoke; + ``` + +**Ожидаемый результат** + +- Таблица `spark_trino_smoke` видна и в Spark, и в Trino. +- Данные совпадают (например, строка `1, from_spark`). + +--- + +## Лаба 3. Партиционирование и эволюция схемы + +**Цель** + +- Показать, как Iceberg работает с партиционированием (фильтрация по partition key, уменьшение объёма чтения). +- Показать, как Iceberg поддерживает эволюцию схемы без сложных миграций. + +**Предусловия** + +- Стенд запущен. +- Желательно выполнить Лабы 1–2, чтобы уже была интуиция про Iceberg и общий каталог. + +**Подготовленные примеры** + +- `src/spark/partitioned_table_demo.sql` — создание партиционированной Iceberg‑таблицы и вставка данных. +- `src/spark/schema_evolution_demo.sql` — демонстрация `ALTER TABLE` и добавления колонок. +- `src/trino/schema_evolution_demo.sql` — чтение той же таблицы с эволюцией схемы из Trino. +- `notebooks/03_partitioning_and_schema_evolution.ipynb` — интерактивный разбор тех же примеров в Jupyter. + +**Вариант A: через Spark SQL и Trino CLI** + +1. Создать партиционированную таблицу в Spark: + + ```bash + docker compose exec spark-master \ + /opt/spark/bin/spark-sql -f /opt/src/spark/partitioned_table_demo.sql + ``` + +2. Проверить в Spark, какие данные записаны по датам: + + ```bash + docker compose exec -it spark-master \ + /opt/spark/bin/spark-sql + ``` + + Внутри интерактивной сессии: + + ```sql + USE lakehouse.default; + SELECT event_date, COUNT(*) AS cnt, SUM(amount) AS total_amount + FROM partition_demo + GROUP BY event_date + ORDER BY event_date; + ``` + +3. Подготовить таблицу с эволюцией схемы: + + ```bash + docker compose exec spark-master \ + /opt/spark/bin/spark-sql -f /opt/src/spark/schema_evolution_demo.sql + ``` + +4. Посмотреть схему и данные в Spark: + + ```bash + docker compose exec -it spark-master \ + /opt/spark/bin/spark-sql + ``` + + Внутри: + + ```sql + USE lakehouse.default; + DESCRIBE TABLE schema_evolution_demo; + SELECT * FROM schema_evolution_demo ORDER BY id; + ``` + +5. Прочитать таблицу с эволюцией схемы из Trino: + + ```bash + docker compose cp src/trino/schema_evolution_demo.sql trino:/tmp/ + docker compose exec trino trino --file /tmp/schema_evolution_demo.sql + ``` + +**Вариант B: через Jupyter‑ноутбук** + +- Открыть `http://localhost:8888` и запустить ноутбук `03_partitioning_and_schema_evolution.ipynb`. +- Ноутбук: + - создаёт SparkSession, подключённый к кластеру; + - выполняет скрипты `partitioned_table_demo.sql` и `schema_evolution_demo.sql`; + - показывает агрегаты по партициям и данные до/после эволюции схемы в интерактивном виде. + +--- + +## Лаба 4. Мини‑ETL поверх Lakehouse (эскиз) + +Идея лабы — собрать end-to-end сценарий: + +- есть сырые данные (CSV/JSON) в S3/MinIO; +- Spark читает raw‑данные, чистит и пишет в Iceberg‑таблицу; +- Trino делает поверх неё аналитику. + +Планируемые компоненты: + +- Папка с примерами сырых данных (например, `examples/raw/` в репозитории, затем загрузка в MinIO). +- `src/spark/etl_raw_to_iceberg.py` — мини ETL, записывающий данные в `lakehouse.default.events` или аналогичную таблицу. +- `src/trino/etl_analytics.sql` — несколько аналитических запросов поверх этой таблицы. + +Детали реализации можно развивать по мере появления новых сценариев. diff --git a/LICENSE b/LICENSE new file mode 100755 index 0000000..5262efb --- /dev/null +++ b/LICENSE @@ -0,0 +1,15 @@ +Creative Commons Attribution 4.0 International (CC BY 4.0) + +This work is licensed under the Creative Commons Attribution 4.0 International License. + +You are free to: +- Share — copy and redistribute the material in any medium or format. +- Adapt — remix, transform, and build upon the material for any purpose, even commercially. + +Under the following terms: +- Attribution — You must give appropriate credit, provide a link to the license, and indicate if changes were made. +- No additional restrictions — You may not apply legal terms or technological measures that legally restrict others from doing anything the license permits. + +The full text of the license is available at: +https://creativecommons.org/licenses/by/4.0/legalcode + diff --git a/README.md b/README.md new file mode 100755 index 0000000..f18642c --- /dev/null +++ b/README.md @@ -0,0 +1,390 @@ +# Учебный Lakehouse-стенд (Trino + Spark + Iceberg + MinIO) + +Учебный стенд для демонстрации **Lakehouse-архитектуры** на одном ноутбуке: + +- объектное хранилище S3-класса (MinIO), +- движок запросов Trino, +- Spark для batch/ETL и интерактивных экспериментов, +- формат таблиц Iceberg, +- JDBC-каталог (PostgreSQL) для метаданных Iceberg в Trino, +- конфиг Spark, заточенный под работу с Iceberg + S3. + +Стенд ориентирован на обучение менти и быструю демонстрацию концепции Lakehouse: разделение **storage / compute / catalog** без лишней обвязки (Hive Metastore, полноценный Hadoop-кластер и т.п.). + +--- + +## Архитектура + +```mermaid +graph LR + subgraph Storage + M[(MinIO
S3-compatible)] + end + + subgraph Catalog + P[(PostgreSQL
Iceberg JDBC catalog)] + end + + subgraph Compute + S[Spark 3.5.1
PySpark / SQL] + T[Trino 478
Iceberg connector] + end + + S <--> M + T <--> M + T <--> P +``` + +**Основные идеи:** + +* **Данные** (Parquet-файлы + служебные каталоги Iceberg) лежат в бакете MinIO (`s3a://lakehouse/warehouse/...`). +* **Метаданные** Iceberg хранятся в PostgreSQL (JDBC-каталог `lakehouse`). +* **Spark и Trino** используют один и тот же JDBC-каталог: таблицы, созданные в Spark, доступны в Trino, и наоборот. + +--- + +## Быстрый старт + +1. Установить Docker и Docker Compose (см. раздел «Требования» ниже). +2. В корне репозитория выполнить: + + ```bash + docker compose build + docker compose up -d + ``` + +3. Открыть основные интерфейсы: + + - Trino Web UI: `http://localhost:8090` + - MinIO Console: `http://localhost:9001` + - JupyterLab (если включён): `http://localhost:8888` + +4. Для пошаговых лабораторных работ см. файл `HOWTO.md`. + +--- + +## Состав репозитория + +* `docker-compose.yml` — описание всех сервисов стенда: + + * Spark master / worker(ы) на базе кастомного образа `spark-iceberg`, + * Trino, + * MinIO (S3-совместимое хранилище), + * PostgreSQL под Iceberg JDBC-каталог, + * (опционально) Jupyter / notebook для PySpark. +* `spark/Dockerfile` — сборка кастомного образа **Spark** с зависимостями: + + * `hadoop-aws`, `aws-java-sdk`, + * библиотека Iceberg нужной версии, + * `pyspark`, `pyarrow` и базовый набор инструментов. +* `jupyter/Dockerfile` — образ JupyterLab на базе Spark-образа. +* `spark/spark-defaults.conf` — конфигурация Spark для работы с: + + * Iceberg-каталогом `lakehouse` (тип `jdbc`, метаданные в Postgres), + * MinIO через `s3a://`, + * расширениями `IcebergSparkSessionExtensions`. +* `src/` — учебные примеры для Spark и Trino: + + * `src/spark/cluster_smoke.py` — проверка, что кластер жив (Spark master/worker). + * `src/spark/iceberg_smoke.py` — Spark создаёт Iceberg-таблицу и читает её. + * `src/spark/iceberg_demo.sql` — пример создания Iceberg-таблицы через Spark SQL. + * `src/trino/iceberg_smoke.sql` — Trino читает таблицу, созданную в Spark. +* `trino/catalog/lakehouse.properties` — конфиг каталога Trino `lakehouse`: + + * коннектор `iceberg`, + * `jdbc`-каталог (PostgreSQL), + * доступ к MinIO как к S3-хранилищу. + +--- + +## Требования + +- Docker и Docker Compose. +- Порты по умолчанию должны быть свободны (см. таблицу «Сервисы и порты» ниже). + +### Сервисы и порты + +| Сервис | Контейнер | Порт (host → container) | Назначение / UI | +|------------------------|-------------------|-------------------------|-------------------------------------| +| Trino | `trino` | `8090 → 8080` | Web UI Trino | +| MinIO API | `minio` | `9000 → 9000` | S3 endpoint | +| MinIO Console | `minio` | `9001 → 9001` | Веб-консоль MinIO | +| PostgreSQL (каталог) | `postgres-iceberg`| `5432 → 5432` | Доступ для psql/DBeaver и т.п. | +| Spark master UI | `spark-master` | `8080 → 8080` | Web UI мастера Spark | +| Spark worker-1 UI | `spark-worker-1` | `8081 → 8081` | Web UI первого воркера | +| Spark worker-2 UI | `spark-worker-2` | `8082 → 8081` | Web UI второго воркера (host 8082) | +| JupyterLab | `jupyter` | `8888 → 8888` | JupyterLab с PySpark | + +--- + +## Сборка и запуск + +Все команды в этом разделе выполняются из корня репозитория. + +### 1. Собрать образы + +```bash +docker compose build +``` + +Будут собраны кастомные образы Spark/Jupyter с зависимостями Iceberg, S3A, JDBC-драйвером Postgres и Python-библиотеками. + +### 2. Поднять стенд + +```bash +docker compose up -d +``` + +Что происходит при старте: + +* MinIO поднимается с root-пользователем/паролем (по умолчанию смотри в `docker-compose.yml`, обычно `minioadmin/minioadmin`), init-контейнер создаёт бакет `lakehouse`. +* PostgreSQL под каталог Iceberg создаёт БД и пользователя (значения — в `docker-compose.yml` / `.env`). +* Trino стартует с каталогом `lakehouse`, описанным в `lakehouse.properties`. +* Spark master/worker получают конфиг из `spark-defaults.conf` (общий JDBC-каталог Iceberg + MinIO). + +Проверить статус: + +```bash +docker compose ps +``` + +### 3. Остановить стенд и очистить данные + +Чтобы остановить все сервисы и удалить данные в MinIO/PostgreSQL (Docker volumes), можно выполнить: + +```bash +docker compose down -v +``` + +--- + +## Доступ к сервисам + +### Trino + +Web UI (координатор): + +```text +http://localhost:8090 +``` + +(точный порт см. в `docker-compose.yml`). + +Подключение **из контейнера trino**: + +```bash +docker exec -it trino trino \ + --server http://localhost:8080 \ + --catalog lakehouse +``` + +Проверка: + +```sql +SHOW CATALOGS; +SHOW SCHEMAS FROM lakehouse; +``` + +### MinIO + +* Консоль: `http://localhost:9001` +* S3 endpoint: `http://localhost:9000` + +Учётные данные — `MINIO_ROOT_USER` / `MINIO_ROOT_PASSWORD` из `docker-compose.yml` или `.env`. + +Создай в MinIO бакет `lakehouse` (если его нет) — данные Iceberg будут храниться именно там. + +### PostgreSQL (каталог Iceberg для Trino) + +Подключение (пример): + +```bash +psql -h localhost -p 5432 -U iceberg -d iceberg +``` + +(имя пользователя, БД и порт уточняются в `docker-compose.yml`). + +Служебные таблицы Iceberg JDBC-каталога (`iceberg_tables`, `iceberg_namespace_properties`) создаёт однократный контейнер `iceberg-catalog-init`. Если база уже запускалась без них, можно переинициализировать вручную: + +```bash +docker compose run --rm iceberg-catalog-init +``` + +### Jupyter / notebooks + +Веб-интерфейс: `http://localhost:8888` (по умолчанию без токена). + +* В контейнере монтируется `./notebooks` в `/opt/work`. +* `./src` доступен read-only в `/opt/src`, переменная `PYTHONPATH=/opt/src` уже установлена — можно импортировать функции из скриптов прямо в ноутбуках. + +--- + +## Конфигурация Trino (каталог `lakehouse`) + +Файл `lakehouse.properties` монтируется в `/etc/trino/catalog/lakehouse.properties`. + +Ключевые параметры: + +```properties +connector.name=iceberg + +# Каталог Iceberg типа JDBC (метаданные в PostgreSQL) +iceberg.catalog.type=jdbc +iceberg.jdbc-catalog.catalog-name=lakehouse +iceberg.jdbc-catalog.driver-class=org.postgresql.Driver +iceberg.jdbc-catalog.connection-url=jdbc:postgresql://postgres-iceberg:5432/iceberg +iceberg.jdbc-catalog.connection-user=iceberg +iceberg.jdbc-catalog.connection-password=iceberg +iceberg.jdbc-catalog.schema-version=V1 + +# Хранилище файлов Iceberg – S3 (MinIO) через hadoop-клиент +iceberg.file-system.type=hadoop +fs.native-s3.enabled=true + +s3.endpoint=http://minio:9000 +s3.region=us-east-1 +s3.path-style-access=true +s3.aws-access-key=minioadmin +s3.aws-secret-key=minioadmin +``` + +**Что это даёт:** + +* Trino хранит *метаданные* Iceberg в PostgreSQL (таблицы каталога, снапшоты, манифесты и т.п.). +* *Файлы данных* лежат в MinIO, в бакете `lakehouse`, к которому Trino ходит по `s3://`/`s3a://` через S3-клиент. + +Пример создания схемы и таблицы из Trino: + +```sql +-- Схема в каталоге lakehouse (метаданные в PostgreSQL) +CREATE SCHEMA lakehouse.default; + +-- Таблица Iceberg с данными в s3://lakehouse/default/test_table/ +CREATE TABLE lakehouse.default.test_table ( + id bigint, + name varchar +); +``` + +--- + +## Конфигурация Spark (spark-defaults.conf) + +`spark-defaults.conf` монтируется в `/opt/spark/conf/spark-defaults.conf` в контейнеры Spark. + +Ключевые моменты: + +```properties +# Iceberg Spark extensions +spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions + +# Каталог Iceberg для Spark (JDBC, метаданные в Postgres) +spark.sql.catalog.lakehouse=org.apache.iceberg.spark.SparkCatalog +spark.sql.catalog.lakehouse.catalog-impl=org.apache.iceberg.jdbc.JdbcCatalog +spark.sql.catalog.lakehouse.uri=jdbc:postgresql://postgres-iceberg:5432/iceberg +spark.sql.catalog.lakehouse.jdbc.user=iceberg +spark.sql.catalog.lakehouse.jdbc.password=iceberg +spark.sql.catalog.lakehouse.jdbc.driver=org.postgresql.Driver +spark.sql.catalog.lakehouse.warehouse=s3a://lakehouse/warehouse +spark.sql.catalog.lakehouse.default-namespace=default + +# S3/MinIO через s3a +spark.hadoop.fs.s3a.endpoint=http://minio:9000 +spark.hadoop.fs.s3a.access.key=minioadmin +spark.hadoop.fs.s3a.secret.key=minioadmin +spark.hadoop.fs.s3a.path.style.access=true +spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem +spark.hadoop.fs.s3a.connection.ssl.enabled=false +``` + +**Важно:** Spark и Trino используют единый JDBC-каталог `lakehouse`: метаданные лежат в Postgres, данные — в MinIO. Таблица, созданная в Spark, видна в Trino без дополнительной настройки. + +Пример создания таблицы из Spark: + +```python +from pyspark.sql import SparkSession + +spark = (SparkSession.builder + .appName("lakehouse-demo") + .getOrCreate()) + +# Каталог lakehouse указан явно +spark.sql(""" + CREATE TABLE lakehouse.default.spark_table ( + id BIGINT, + name STRING + ) + USING iceberg +""") + +spark.sql("INSERT INTO lakehouse.default.spark_table VALUES (1, 'Alice'), (2, 'Bob')") +``` + +Файлы окажутся в `s3a://lakehouse/warehouse/default/spark_table/`. + +--- + +## Типовой учебный сценарий + +1. **Поднять стенд** (`docker compose build`, затем `docker compose up -d`). +2. **Создать бакет `lakehouse`** в MinIO Console (если ещё нет). +3. **Создать таблицу из Spark**, записать туда данные, показать дерево файлов Iceberg в MinIO (data/manifest/metadata). +4. **Прочитать ту же таблицу из Trino** (проверка общего каталога). +5. **Создать таблицу из Trino** и прочитать её из Spark. +6. Обсудить архитектуру: Postgres хранит метаданные Iceberg, MinIO — данные, Spark/Trino — compute. + +--- + +## Smoke-тесты Spark → Trino + +Быстрая проверка, что таблица, созданная в Spark, читается в Trino через общий JDBC-каталог. + +1. Убедиться, что стенд запущен: `docker compose up -d`. +2. Скопировать скрипты внутрь контейнеров: + + ```bash + docker compose cp src/spark/iceberg_smoke.py spark-master:/tmp/ + docker compose cp src/trino/iceberg_smoke.sql trino:/tmp/ + ``` + +3. Выполнить smoke из Spark: + + ```bash + docker compose exec spark-master /opt/spark/bin/spark-submit /tmp/iceberg_smoke.py + ``` + + Скрипт создаст `lakehouse.default.spark_trino_smoke`, вставит строку `1, from_spark` и прочитает её. + +4. Прочитать ту же таблицу из Trino: + + ```bash + docker compose exec trino trino --file /tmp/iceberg_smoke.sql + ``` + + В выводе должны быть строки `spark_trino_smoke` в списке таблиц и `1, from_spark` в результате выборки. + +5. Очистка (опционально): + + ```bash + docker compose exec spark-master /opt/spark/bin/spark-sql -e "DROP TABLE IF EXISTS lakehouse.default.spark_trino_smoke" + ``` + +--- + +## Дальнейшее развитие + +Планируемые/возможные расширения: + +* Подключить **Airflow** и запускать Spark-job’ы поверх этого же Lakehouse. +* Вынести настройки (`MINIO_ROOT_USER`, `ICEBERG_*`, `POSTGRES_*`) в `.env` с шаблоном для студентов. +* Добавить отдельные каталоги Trino (например, `hive`, `tpch`) для демонстрации федеративных запросов. +* Добавить пример интеграции с BI-инструментом (DBeaver/Metabase/Superset) поверх Trino. + +Для пошаговых учебных сценариев (лабораторных работ) см. файл `HOWTO.md`. + +--- + +## Лицензия + +Материалы этого репозитория лицензированы на условиях Creative Commons Attribution 4.0 International (CC BY 4.0). +См. файл `LICENSE` или . diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100755 index 0000000..39451c7 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,209 @@ +# version: "3.9" +services: + postgres: + image: postgres:16 + container_name: postgres-iceberg + environment: + POSTGRES_DB: iceberg # БД под Iceberg catalog + POSTGRES_USER: iceberg + POSTGRES_PASSWORD: iceberg + ports: + - "5432:5432" # опционально, чтобы заходить с хоста через DBeaver/psql + networks: + - spark-net + volumes: + - postgres-data:/var/lib/postgresql/data + healthcheck: + test: ["CMD-SHELL", "pg_isready -U iceberg -d iceberg || exit 1"] + interval: 10s + timeout: 5s + retries: 5 + + iceberg-catalog-init: + image: postgres:16 + container_name: iceberg-catalog-init + depends_on: + - postgres + networks: + - spark-net + environment: + - PGPASSWORD=iceberg + entrypoint: | + /bin/sh -c "set -e + echo 'Waiting for Postgres (iceberg catalog)...' + until pg_isready -h postgres-iceberg -U iceberg -d iceberg; do + sleep 1 + done + cat <<'SQL' | psql -h postgres-iceberg -U iceberg -d iceberg + CREATE TABLE IF NOT EXISTS iceberg_tables ( + catalog_name VARCHAR(255) NOT NULL, + table_namespace VARCHAR(255) NOT NULL, + table_name VARCHAR(255) NOT NULL, + metadata_location VARCHAR(1000), + previous_metadata_location VARCHAR(1000), + iceberg_type VARCHAR(5), + PRIMARY KEY (catalog_name, table_namespace, table_name) + ); + CREATE TABLE IF NOT EXISTS iceberg_namespace_properties ( + catalog_name VARCHAR(255) NOT NULL, + namespace VARCHAR(255) NOT NULL, + property_key VARCHAR(255), + property_value VARCHAR(1000), + PRIMARY KEY (catalog_name, namespace, property_key) + ); + SQL + echo 'Iceberg JDBC catalog tables are ready' + " + + spark-master: + build: + context: ./spark + image: de-spark-iceberg:3.5.1 + container_name: spark-master + hostname: spark-master + depends_on: + - postgres + networks: + - spark-net + environment: + - SPARK_NO_DAEMONIZE=true + - AWS_ACCESS_KEY_ID=minioadmin + - AWS_SECRET_ACCESS_KEY=minioadmin + command: ["/opt/spark/sbin/start-master.sh"] + ports: + - "7077:7077" # Spark master + - "8080:8080" # Web UI мастера + volumes: + - ./spark/spark-defaults.conf:/opt/spark/conf/spark-defaults.conf:ro + + spark-worker-1: + image: de-spark-iceberg:3.5.1 + # build: + # context: ./spark + container_name: spark-worker-1 + hostname: spark-worker-1 + networks: + - spark-net + environment: + - SPARK_NO_DAEMONIZE=true + - SPARK_WORKER_CORES=2 + - SPARK_WORKER_MEMORY=2g + - AWS_ACCESS_KEY_ID=minioadmin + - AWS_SECRET_ACCESS_KEY=minioadmin + command: ["/opt/spark/sbin/start-worker.sh", "spark://spark-master:7077"] + depends_on: + - spark-master + ports: + - "8081:8081" # Web UI первого воркера + volumes: + - ./spark/spark-defaults.conf:/opt/spark/conf/spark-defaults.conf:ro + + spark-worker-2: + image: de-spark-iceberg:3.5.1 + # build: + # context: ./spark + container_name: spark-worker-2 + hostname: spark-worker-2 + networks: + - spark-net + environment: + - SPARK_NO_DAEMONIZE=true + - SPARK_WORKER_CORES=2 + - SPARK_WORKER_MEMORY=2g + - AWS_ACCESS_KEY_ID=minioadmin + - AWS_SECRET_ACCESS_KEY=minioadmin + command: ["/opt/spark/sbin/start-worker.sh", "spark://spark-master:7077"] + depends_on: + - spark-master + ports: + - "8082:8081" # Web UI второго воркера (host 8082) + volumes: + - ./spark/spark-defaults.conf:/opt/spark/conf/spark-defaults.conf:ro + + minio: + image: minio/minio:RELEASE.2024-10-02T17-50-41Z + container_name: minio + hostname: minio + networks: + - spark-net + command: server /data --console-address ":9001" + environment: + - MINIO_ROOT_USER=minioadmin + - MINIO_ROOT_PASSWORD=minioadmin + ports: + - "9000:9000" + - "9001:9001" + volumes: + - minio-data:/data + + minio-init: + image: minio/mc:RELEASE.2024-06-24T19-40-33Z-cpuv1 + container_name: minio-init + hostname: minio-init + networks: + - spark-net + depends_on: + - minio + entrypoint: > + /bin/sh -c " + echo 'Waiting for MinIO...'; + until mc alias set minio http://minio:9000 minioadmin minioadmin 2>/dev/null; do + sleep 1; + done; + echo 'MinIO is up, creating bucket...'; + mc mb -p minio/lakehouse || true; + echo 'Bucket lakehouse is ready'; + exit 0; + " + + + jupyter: + build: + context: ./jupyter + args: + NB_UID: ${NB_UID:-1000} + NB_GID: ${NB_GID:-100} + # image: jupyter-spark-iceberg:3.5.1 + container_name: jupyter + hostname: jupyter + depends_on: + - spark-master + - minio + networks: + - spark-net + environment: + - AWS_ACCESS_KEY_ID=minioadmin + - AWS_SECRET_ACCESS_KEY=minioadmin + - PYTHONPATH=/opt/src + ports: + - "8888:8888" + volumes: + - ./spark/spark-defaults.conf:/opt/spark/conf/spark-defaults.conf:ro + - ./notebooks:/opt/work + - ./src:/opt/src:ro + + trino: + image: trinodb/trino:478 + container_name: trino + hostname: trino + networks: + - spark-net + depends_on: + - minio + - postgres + - iceberg-catalog-init + environment: + - AWS_ACCESS_KEY_ID=minioadmin + - AWS_SECRET_ACCESS_KEY=minioadmin + - AWS_REGION=us-east-1 + ports: + - "8090:8080" # Web UI Trino + volumes: + - ./trino/catalog:/etc/trino/catalog:ro + +networks: + spark-net: + +volumes: + postgres-data: + minio-data: diff --git a/jupyter/Dockerfile b/jupyter/Dockerfile new file mode 100755 index 0000000..daf1928 --- /dev/null +++ b/jupyter/Dockerfile @@ -0,0 +1,27 @@ +# Берем за основу контейнер, приготовленный для spark в spark/Dockerfile +FROM de-spark-iceberg:3.5.1 + +USER root +ENV DEBIAN_FRONTEND=noninteractive + +# Только JupyterLab, без дополнительных пакетов +RUN pip3 install --no-cache-dir "jupyterlab==4.2.5" + +# Создаём непривилегированного пользователя +ARG NB_USER=jovyan +ARG NB_UID=1000 + +RUN useradd --uid ${NB_UID} -m -s /bin/bash ${NB_USER} \ + && mkdir -p /opt/work \ + && chown -R ${NB_UID}:${NB_UID} /opt/work + +USER ${NB_UID} +WORKDIR /opt/work +EXPOSE 8888 + +CMD ["jupyter", "lab", \ + "--ip=0.0.0.0", \ + "--port=8888", \ + "--no-browser", \ + "--NotebookApp.token=", \ + "--NotebookApp.password="] diff --git a/notebooks/03_partitioning_and_schema_evolution.ipynb b/notebooks/03_partitioning_and_schema_evolution.ipynb new file mode 100755 index 0000000..52e5c05 --- /dev/null +++ b/notebooks/03_partitioning_and_schema_evolution.ipynb @@ -0,0 +1,163 @@ +{ + "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 +} diff --git a/notebooks/spark-basic-test.ipynb b/notebooks/spark-basic-test.ipynb new file mode 100755 index 0000000..623fb20 --- /dev/null +++ b/notebooks/spark-basic-test.ipynb @@ -0,0 +1,86 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 1, + "id": "ca5e5822-18cb-405e-90f3-a0f2f8ad2260", + "metadata": {}, + "outputs": [ + { + "name": "stderr", + "output_type": "stream", + "text": [ + "Setting default log level to \"WARN\".\n", + "To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).\n", + "25/12/03 09:08:48 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable\n", + "25/12/03 09:08:52 WARN MetricsConfig: Cannot locate configuration: tried hadoop-metrics2-s3a-file-system.properties,hadoop-metrics2.properties\n", + " \r" + ] + }, + { + "name": "stdout", + "output_type": "stream", + "text": [ + "+---+---+\n", + "| id|txt|\n", + "+---+---+\n", + "| 1| ok|\n", + "+---+---+\n", + "\n" + ] + } + ], + "source": [ + "from pyspark.sql import SparkSession\n", + "\n", + "spark = (\n", + " SparkSession.builder\n", + " .appName(\"check\")\n", + " .master(\"spark://spark-master:7077\")\n", + " .getOrCreate()\n", + ")\n", + "\n", + "spark.sql(\"\"\"\n", + " CREATE TABLE IF NOT EXISTS lakehouse.default.demo_fix (\n", + " id BIGINT,\n", + " txt STRING\n", + " )\n", + " USING iceberg\n", + "\"\"\")\n", + "\n", + "spark.sql(\"INSERT INTO lakehouse.default.demo_fix VALUES (1, 'ok')\")\n", + "\n", + "spark.sql(\"SELECT * FROM lakehouse.default.demo_fix\").show()\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "1e125075-feab-4974-baa2-f787e41b46fd", + "metadata": {}, + "outputs": [], + "source": [] + } + ], + "metadata": { + "kernelspec": { + "display_name": "Python 3 (ipykernel)", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.10.12" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} diff --git a/spark/Dockerfile b/spark/Dockerfile new file mode 100755 index 0000000..0807f16 --- /dev/null +++ b/spark/Dockerfile @@ -0,0 +1,47 @@ +FROM spark:3.5.1-scala2.12-java17-python3-ubuntu + +USER root +ENV DEBIAN_FRONTEND=noninteractive + +# Версии зависимостей в одном месте +ARG HADOOP_AWS_VERSION=3.3.4 +ARG AWS_SDK_VERSION=1.12.262 +ARG ICEBERG_VERSION=1.10.0 +ARG POSTGRES_VERSION=42.7.3 +ARG PYSPARK_VERSION=3.5.1 + +# Базовый Python/утилиты +RUN apt-get update \ + && apt-get install -y python3-pip curl wget nano mc tmux\ + && rm -rf /var/lib/apt/lists/* + +# JAR'ы для S3A и Iceberg +RUN mkdir -p /opt/spark/jars \ + && curl -L "https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/${HADOOP_AWS_VERSION}/hadoop-aws-${HADOOP_AWS_VERSION}.jar" \ + -o "/opt/spark/jars/hadoop-aws-${HADOOP_AWS_VERSION}.jar" \ + && curl -L "https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/${AWS_SDK_VERSION}/aws-java-sdk-bundle-${AWS_SDK_VERSION}.jar" \ + -o "/opt/spark/jars/aws-java-sdk-bundle-${AWS_SDK_VERSION}.jar" \ + && curl -L "https://repo1.maven.org/maven2/org/apache/iceberg/iceberg-spark-runtime-3.5_2.12/${ICEBERG_VERSION}/iceberg-spark-runtime-3.5_2.12-${ICEBERG_VERSION}.jar" \ + -o "/opt/spark/jars/iceberg-spark-runtime-3.5_2.12-${ICEBERG_VERSION}.jar" \ + && curl -L "https://repo1.maven.org/maven2/org/postgresql/postgresql/${POSTGRES_VERSION}/postgresql-${POSTGRES_VERSION}.jar" \ + -o "/opt/spark/jars/postgresql-${POSTGRES_VERSION}.jar" + +# Общие Python-пакеты для драйвера и executors +RUN pip3 install --no-cache-dir \ + "pyspark==${PYSPARK_VERSION}" \ + pandas \ + numpy \ + pyarrow + +ENV PYSPARK_PYTHON=python3 +ENV PYSPARK_DRIVER_PYTHON=python3 + +# Конфиг Spark (S3 + Iceberg + каталог lakehouse) +RUN mkdir -p /opt/spark/conf +# Конфиги подключим уже в docker-compose +# COPY spark-defaults.conf /opt/spark/conf/spark-defaults.conf + +# Возвращаемся к непривилегированному UID как в официальном образе +# (комментарий отдельно, чтобы не ломать USER) +# тот же пользователь, что в официальном spark-образе +USER 185 diff --git a/spark/spark-defaults.conf b/spark/spark-defaults.conf new file mode 100755 index 0000000..b626005 --- /dev/null +++ b/spark/spark-defaults.conf @@ -0,0 +1,26 @@ +# spark/spark-defaults.conf + +# Включаем Iceberg extensions +spark.sql.extensions org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions + +# JDBC-каталог "lakehouse" (метаданные в Postgres, данные в MinIO) +spark.sql.catalog.lakehouse org.apache.iceberg.spark.SparkCatalog +spark.sql.catalog.lakehouse.catalog-impl org.apache.iceberg.jdbc.JdbcCatalog +spark.sql.catalog.lakehouse.uri jdbc:postgresql://postgres-iceberg:5432/iceberg +spark.sql.catalog.lakehouse.jdbc.user iceberg +spark.sql.catalog.lakehouse.jdbc.password iceberg +spark.sql.catalog.lakehouse.jdbc.driver org.postgresql.Driver +spark.sql.catalog.lakehouse.warehouse s3a://lakehouse/warehouse + +# Настройки S3A для MinIO +spark.hadoop.fs.s3a.impl org.apache.hadoop.fs.s3a.S3AFileSystem +spark.hadoop.fs.s3a.endpoint http://minio:9000 +spark.hadoop.fs.s3a.path.style.access true +spark.hadoop.fs.s3a.connection.ssl.enabled false + +# Креды MinIO (для учебного стенда ок в явном виде) +spark.hadoop.fs.s3a.access.key minioadmin +spark.hadoop.fs.s3a.secret.key minioadmin + +# Чуть-чуть логики по умолчанию +spark.sql.catalog.lakehouse.default-namespace default diff --git a/src/spark/cluster_smoke.py b/src/spark/cluster_smoke.py new file mode 100755 index 0000000..a222785 --- /dev/null +++ b/src/spark/cluster_smoke.py @@ -0,0 +1,12 @@ +from pyspark.sql import SparkSession + +spark = ( + SparkSession.builder + .appName("test-cluster") + .master("spark://spark-master:7077") # важно: не local[*] + .getOrCreate() +) + +spark.range(0, 1000000).groupBy().sum().show() + +spark.stop() diff --git a/src/spark/iceberg_demo.sql b/src/spark/iceberg_demo.sql new file mode 100755 index 0000000..326ddc1 --- /dev/null +++ b/src/spark/iceberg_demo.sql @@ -0,0 +1,10 @@ +-- создаём таблицу в каталоге lakehouse, в MinIO +CREATE TABLE lakehouse.default.demo_tbl ( + id BIGINT, + txt STRING +) +USING iceberg; + +INSERT INTO lakehouse.default.demo_tbl VALUES (1, 'hello'), (2, 'world'); + +SELECT * FROM lakehouse.default.demo_tbl; diff --git a/src/spark/iceberg_smoke.py b/src/spark/iceberg_smoke.py new file mode 100755 index 0000000..e7c2b86 --- /dev/null +++ b/src/spark/iceberg_smoke.py @@ -0,0 +1,24 @@ +#!/usr/bin/env python3 +from pyspark.sql import SparkSession + + +TABLE = "lakehouse.default.spark_trino_smoke" + + +def main() -> None: + spark = SparkSession.builder.appName("iceberg-smoke").getOrCreate() + + spark.sql("CREATE NAMESPACE IF NOT EXISTS lakehouse.default") + spark.sql(f"DROP TABLE IF EXISTS {TABLE}") + spark.sql(f"CREATE TABLE {TABLE} (id INT, payload STRING) USING iceberg") + spark.sql(f"INSERT INTO {TABLE} VALUES (1, 'from_spark')") + + result = spark.sql(f"SELECT * FROM {TABLE} ORDER BY id") + result.show(truncate=False) + print(f"✔ Spark записал и прочитал таблицу {TABLE}") + + spark.stop() + + +if __name__ == "__main__": + main() diff --git a/src/spark/partitioned_table_demo.sql b/src/spark/partitioned_table_demo.sql new file mode 100755 index 0000000..76db5ff --- /dev/null +++ b/src/spark/partitioned_table_demo.sql @@ -0,0 +1,20 @@ +-- Партиционированная Iceberg-таблица lakehouse.default.partition_demo + +CREATE NAMESPACE IF NOT EXISTS lakehouse.default; + +DROP TABLE IF EXISTS lakehouse.default.partition_demo; + +CREATE TABLE lakehouse.default.partition_demo ( + user_id BIGINT, + event_date DATE, + amount DOUBLE +) +USING iceberg +PARTITIONED BY (event_date); + +INSERT INTO lakehouse.default.partition_demo VALUES + (1, '2024-01-01', 10.0), + (2, '2024-01-01', 20.0), + (3, '2024-02-01', 30.0), + (4, '2024-02-01', 40.0); + diff --git a/src/spark/schema_evolution_demo.sql b/src/spark/schema_evolution_demo.sql new file mode 100755 index 0000000..8f740a8 --- /dev/null +++ b/src/spark/schema_evolution_demo.sql @@ -0,0 +1,23 @@ +-- Эволюция схемы Iceberg-таблицы lakehouse.default.schema_evolution_demo + +CREATE NAMESPACE IF NOT EXISTS lakehouse.default; + +DROP TABLE IF EXISTS lakehouse.default.schema_evolution_demo; + +CREATE TABLE lakehouse.default.schema_evolution_demo ( + id BIGINT, + value STRING +) +USING iceberg; + +INSERT INTO lakehouse.default.schema_evolution_demo VALUES + (1, 'old_row_1'), + (2, 'old_row_2'); + +ALTER TABLE lakehouse.default.schema_evolution_demo + ADD COLUMN metadata STRING; + +INSERT INTO lakehouse.default.schema_evolution_demo VALUES + (3, 'new_row_1', 'extra_info'), + (4, 'new_row_2', 'extra_info'); + diff --git a/src/trino/iceberg_smoke.sql b/src/trino/iceberg_smoke.sql new file mode 100755 index 0000000..c918bdb --- /dev/null +++ b/src/trino/iceberg_smoke.sql @@ -0,0 +1,5 @@ +-- Проверка чтения данных, созданных в Spark +USE lakehouse.default; + +SHOW TABLES; +SELECT * FROM spark_trino_smoke; diff --git a/src/trino/schema_evolution_demo.sql b/src/trino/schema_evolution_demo.sql new file mode 100755 index 0000000..2b9c7ac --- /dev/null +++ b/src/trino/schema_evolution_demo.sql @@ -0,0 +1,13 @@ +-- Проверка эволюции схемы в Trino для таблицы, +-- созданной в Spark (lakehouse.default.schema_evolution_demo). + +USE lakehouse.default; + +SHOW TABLES LIKE 'schema_evolution_demo'; + +DESCRIBE schema_evolution_demo; + +SELECT * +FROM schema_evolution_demo +ORDER BY id; + diff --git a/trino/catalog/lakehouse.properties b/trino/catalog/lakehouse.properties new file mode 100755 index 0000000..05e3afa --- /dev/null +++ b/trino/catalog/lakehouse.properties @@ -0,0 +1,31 @@ +connector.name=iceberg + +# 1. Тип каталога — JDBC (метаданные в Postgres) +iceberg.catalog.type=jdbc + +# Подключение к Postgres с метаданными Iceberg +iceberg.jdbc-catalog.connection-url=jdbc:postgresql://postgres-iceberg:5432/iceberg +iceberg.jdbc-catalog.connection-user=iceberg +iceberg.jdbc-catalog.connection-password=iceberg + +# Обязательные поля JDBC-каталога +iceberg.jdbc-catalog.driver-class=org.postgresql.Driver +iceberg.jdbc-catalog.catalog-name=lakehouse +iceberg.jdbc-catalog.schema-version=V1 + +# 2. Warehouse целиком в S3 (MinIO) +# Здесь "lakehouse" — имя S3-бакета, "warehouse" — префикс внутри него. +iceberg.jdbc-catalog.default-warehouse-dir=s3://lakehouse/warehouse + +# 3. Включаем новый native S3 filesystem +fs.hadoop.enabled=false +fs.native-s3.enabled=true + +# 4. Настройки доступа к MinIO +s3.endpoint=http://minio:9000 +s3.region=us-east-1 +s3.path-style-access=true + +# Статические ключи для MinIO (учебный стенд) +s3.aws-access-key=minioadmin +s3.aws-secret-key=minioadmin