первый commit

This commit is contained in:
2025-12-03 16:35:50 +03:00
commit 41bbdecba5
19 changed files with 1478 additions and 0 deletions
Executable
+23
View File
@@ -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
Executable
+56
View File
@@ -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 <service>` 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.
Executable
+298
View File
@@ -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` — несколько аналитических запросов поверх этой таблицы.
Детали реализации можно развивать по мере появления новых сценариев.
Executable
+15
View File
@@ -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
Executable
+390
View File
@@ -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<br/>S3-compatible)]
end
subgraph Catalog
P[(PostgreSQL<br/>Iceberg JDBC catalog)]
end
subgraph Compute
S[Spark 3.5.1<br/>PySpark / SQL]
T[Trino 478<br/>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` или <https://creativecommons.org/licenses/by/4.0/>.
+209
View File
@@ -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:
+27
View File
@@ -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="]
+163
View File
@@ -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
}
+86
View File
@@ -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
}
+47
View File
@@ -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
+26
View File
@@ -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
+12
View File
@@ -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()
+10
View File
@@ -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;
+24
View File
@@ -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()
+20
View File
@@ -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);
+23
View File
@@ -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');
+5
View File
@@ -0,0 +1,5 @@
-- Проверка чтения данных, созданных в Spark
USE lakehouse.default;
SHOW TABLES;
SELECT * FROM spark_trino_smoke;
+13
View File
@@ -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;
+31
View File
@@ -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