diff --git a/AGENTS.md b/AGENTS.md index 7b47720..745e156 100755 --- a/AGENTS.md +++ b/AGENTS.md @@ -10,6 +10,7 @@ This repository contains a small teaching Lakehouse stack (Spark + Trino + Icebe - `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`). +- `plans/`: internal living docs for implementation plans. Keep new examples in `src/` or `notebooks/`, and avoid mixing configuration and code. @@ -37,7 +38,7 @@ 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. +- opening `notebooks/01_environment_and_smoke_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. diff --git a/README.md b/README.md index 85963df..5190000 100755 --- a/README.md +++ b/README.md @@ -46,6 +46,8 @@ graph LR ## Быстрый старт +Подробный маршрут старта вынесен в [START_HERE.md](./START_HERE.md). Если нужен только краткий запуск, достаточно: + 1. Установить Docker и Docker Compose (см. раздел «Требования» ниже). 2. В корне репозитория выполнить: @@ -56,11 +58,13 @@ graph LR 3. Открыть основные интерфейсы: + - Spark Master UI: `http://localhost:8080` - Trino Web UI: `http://localhost:8090` - MinIO Console: `http://localhost:9001` - JupyterLab (если включён): `http://localhost:8888` -4. Актуальная структура курса описана в `docs/course_program.md`. +4. Для входа в курс и Модуль 1 открыть `START_HERE.md`. +5. Актуальная структура курса описана в `docs/course_program.md`. --- @@ -79,6 +83,8 @@ graph LR * библиотека Iceberg нужной версии, * `pyspark`, `pyarrow` и базовый набор инструментов. * `jupyter/Dockerfile` — образ JupyterLab на базе Spark-образа. +* `START_HERE.md` — канонический маршрут входа в курс до первого ноутбука. +* `plans/` — внутренние living docs с планами работ по модулям курса. * `spark/spark-defaults.conf` — конфигурация Spark для работы с: * Iceberg-каталогом `lakehouse` (тип `jdbc`, метаданные в Postgres), @@ -217,6 +223,7 @@ docker compose run --rm iceberg-catalog-init * В контейнере монтируется `./notebooks` в `/opt/work`. * `./src` доступен read-only в `/opt/src`, переменная `PYTHONPATH=/opt/src` уже установлена — можно импортировать функции из скриптов прямо в ноутбуках. +* Первый канонический ноутбук курса: `notebooks/01_environment_and_smoke_test.ipynb`. --- diff --git a/START_HERE.md b/START_HERE.md new file mode 100644 index 0000000..efebd2d --- /dev/null +++ b/START_HERE.md @@ -0,0 +1,153 @@ +# START HERE + +Этот документ нужен для первого входа в курс и прохождения Модуля 1: поднять стенд, проверить сервисы, открыть интерфейсы и выполнить базовый smoke test. + +## Маршрут прохождения + +1. Проверить prerequisites. +2. Собрать и поднять стенд. +3. Убедиться, что контейнеры живы. +4. Открыть основные UI. +5. Зайти в Jupyter и выполнить `notebooks/01_environment_and_smoke_test.ipynb`. +6. При проблемах использовать логи и шаги диагностики из этого документа. + +## Prerequisites + +Нужно заранее установить: + +- Docker; +- Docker Compose; +- современный браузер для UI; +- свободные порты на хосте. + +Рекомендуемые ресурсы хоста: + +- не менее `4 vCPU`; +- не менее `10 GB RAM`, иначе `Spark`, `Trino` и `Jupyter` могут стартовать нестабильно; +- хотя бы `8-10 GB` свободного места под образы и контейнеры. + +## Порты стенда + +| Сервис | Адрес | Зачем нужен | +| --- | --- | --- | +| Spark Master UI | `http://localhost:8080` | Проверка мастера Spark и подключённых worker-ов | +| Spark Worker 1 UI | `http://localhost:8081` | Проверка первого worker-а | +| Spark Worker 2 UI | `http://localhost:8082` | Проверка второго worker-а | +| Trino UI | `http://localhost:8090` | Проверка координатора Trino | +| MinIO API | `http://localhost:9000` | S3-compatible endpoint | +| MinIO Console | `http://localhost:9001` | Просмотр бакетов и файлов | +| JupyterLab | `http://localhost:8888` | Основная точка входа в практику | +| PostgreSQL | `localhost:5432` | JDBC-каталог Iceberg, нужен для диагностики | + +## Быстрый запуск + +Все команды выполняются из корня репозитория. + +### 1. Собрать образы + +```bash +docker compose build +``` + +### 2. Поднять стенд + +```bash +docker compose up -d +``` + +### 3. Проверить статус контейнеров + +```bash +docker compose ps +``` + +Ожидаемое состояние: + +- сервисы `spark-master`, `spark-worker-1`, `spark-worker-2`, `minio`, `postgres-iceberg`, `trino`, `jupyter` находятся в состоянии `Up`; +- однократные init-контейнеры вроде `minio-init` и `iceberg-catalog-init` могут завершиться после успешной инициализации. + +## Куда заходить после старта + +Открой в браузере: + +- `http://localhost:8080` для `Spark Master UI`; +- `http://localhost:8090` для `Trino UI`; +- `http://localhost:9001` для `MinIO Console`; +- `http://localhost:8888` для `JupyterLab`. + +Для входа в `MinIO Console` используй: + +- логин: `minioadmin`; +- пароль: `minioadmin`. + +Если интерфейсы открываются, переходи в Jupyter и запускай: + +```text +notebooks/01_environment_and_smoke_test.ipynb +``` + +## Роли сервисов в стенде + +| Сервис | Роль в модуле | +| --- | --- | +| `MinIO` | `storage`: объектное хранилище для данных и файлов Iceberg | +| `PostgreSQL` | `catalog`: хранит метаданные JDBC-каталога Iceberg | +| `Spark` | `compute`: выполняет PySpark и Spark SQL задания | +| `Trino` | `compute`: читает те же таблицы через общий каталог | +| `Jupyter` | Точка входа в учебные ноутбуки | + +## Как работать в Jupyter + +- `./notebooks` смонтирован в контейнер как `/opt/work`; +- `./src` смонтирован read-only как `/opt/src`; +- `PYTHONPATH=/opt/src`, поэтому helper-скрипты из `src/` доступны для импорта в ноутбуках; +- первый ноутбук курса: `01_environment_and_smoke_test.ipynb`. + +## Базовая диагностика + +### Посмотреть список контейнеров + +```bash +docker compose ps +``` + +### Посмотреть логи конкретного сервиса + +```bash +docker compose logs -f spark-master +docker compose logs -f trino +docker compose logs -f minio +docker compose logs -f jupyter +``` + +### Типовые первые проверки + +- `Spark UI` не открывается: проверь `docker compose ps` и логи `spark-master`. +- `Trino UI` не открывается: проверь `docker compose ps` и логи `trino`. +- `JupyterLab` не открывается: проверь `docker compose ps` и логи `jupyter`. +- `MinIO Console` не открывается: проверь `docker compose ps` и логи `minio`. +- smoke test падает из ноутбука: сначала убедись, что `spark-master` и worker-ы видны в `Spark UI`. + +## Restart и reset + +### Мягкий перезапуск стенда + +```bash +docker compose down +docker compose up -d +``` + +### Полный reset с удалением данных + +```bash +docker compose down -v +docker compose up -d +``` + +Используй полный reset, если хочешь пройти практику заново с чистого состояния. + +## Что делать дальше + +- пройти `notebooks/01_environment_and_smoke_test.ipynb`; +- свериться с программой курса в `docs/course_program.md`; +- после прохождения Модуля 1 переходить к следующим учебным материалам. diff --git a/docs/archive/legacy_howto.md b/docs/archive/legacy_howto.md index ed076b5..be62ab9 100644 --- a/docs/archive/legacy_howto.md +++ b/docs/archive/legacy_howto.md @@ -234,7 +234,7 @@ - `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. +- интерактивный ноутбук для этой legacy-лабы удалён из актуального репозитория, чтобы не смешивать его с основным треком курса. **Вариант A: через Spark SQL и Trino CLI** @@ -291,13 +291,13 @@ docker compose exec trino trino --file /tmp/schema_evolution_demo.sql ``` -**Вариант B: через Jupyter-ноутбук** +**Вариант B: через Jupyter вручную** -- Открыть `http://localhost:8888` и запустить ноутбук `03_partitioning_and_schema_evolution.ipynb`. -- Ноутбук: - - создаёт SparkSession, подключённый к кластеру; - - выполняет скрипты `partitioned_table_demo.sql` и `schema_evolution_demo.sql`; - - показывает агрегаты по партициям и данные до/после эволюции схемы в интерактивном виде. +- Открыть `http://localhost:8888`, если нужен самостоятельный разбор SQL-примеров в Jupyter. +- Повторить те же шаги вручную в кодовых ячейках: + - создать SparkSession, подключённый к кластеру; + - выполнить скрипты `partitioned_table_demo.sql` и `schema_evolution_demo.sql`; + - посмотреть агрегаты по партициям и данные до/после эволюции схемы в интерактивном виде. --- diff --git a/jupyter/Dockerfile b/jupyter/Dockerfile index daf1928..cacf393 100755 --- a/jupyter/Dockerfile +++ b/jupyter/Dockerfile @@ -10,12 +10,14 @@ RUN pip3 install --no-cache-dir "jupyterlab==4.2.5" # Создаём непривилегированного пользователя ARG NB_USER=jovyan ARG NB_UID=1000 +ARG NB_GID=100 -RUN useradd --uid ${NB_UID} -m -s /bin/bash ${NB_USER} \ +RUN if ! getent group ${NB_GID} >/dev/null; then groupadd --gid ${NB_GID} ${NB_USER}; fi \ + && useradd --uid ${NB_UID} --gid ${NB_GID} -m -s /bin/bash ${NB_USER} \ && mkdir -p /opt/work \ - && chown -R ${NB_UID}:${NB_UID} /opt/work + && chown -R ${NB_UID}:${NB_GID} /opt/work -USER ${NB_UID} +USER ${NB_UID}:${NB_GID} WORKDIR /opt/work EXPOSE 8888 diff --git a/notebooks/01_environment_and_smoke_test.ipynb b/notebooks/01_environment_and_smoke_test.ipynb new file mode 100644 index 0000000..024d0fd --- /dev/null +++ b/notebooks/01_environment_and_smoke_test.ipynb @@ -0,0 +1,213 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "id": "module1-intro", + "metadata": {}, + "source": [ + "# Модуль 1. Вход в стенд и базовая диагностика\n", + "\n", + "Этот ноутбук является канонической практикой Модуля 1.\n", + "\n", + "На выходе ты должен:\n", + "\n", + "- понимать, что стенд уже поднят на хосте;\n", + "- подключиться к Spark-кластеру из Jupyter;\n", + "- выполнить базовый smoke test;\n", + "- сопоставить сервисы стенда с их ролями;\n", + "- знать первый маршрут диагностики через `docker compose ps`, логи и UI.\n" + ] + }, + { + "cell_type": "markdown", + "id": "module1-before-start", + "metadata": {}, + "source": [ + "## Перед стартом\n", + "\n", + "Перед запуском этого ноутбука пройди шаги из `START_HERE.md`.\n", + "\n", + "Важно: `docker compose build`, `docker compose up -d`, `docker compose ps` и просмотр логов выполняются на хосте, а не внутри Jupyter.\n", + "\n", + "Ожидаемое состояние перед началом:\n", + "\n", + "- контейнеры `spark-master`, `spark-worker-1`, `spark-worker-2`, `minio`, `postgres-iceberg`, `trino`, `jupyter` уже подняты;\n", + "- открываются `Spark UI`, `Trino UI`, `MinIO Console` и `JupyterLab`;\n", + "- ты работаешь в Jupyter внутри контейнера `jupyter`.\n" + ] + }, + { + "cell_type": "markdown", + "id": "module1-service-map", + "metadata": {}, + "source": [ + "## Карта сервисов\n", + "\n", + "| Сервис | Роль |\n", + "| --- | --- |\n", + "| `MinIO` | `storage`: объектное хранилище для данных и файлов Iceberg |\n", + "| `PostgreSQL` | `catalog`: хранит метаданные JDBC-каталога Iceberg |\n", + "| `Spark` | `compute`: выполняет PySpark и Spark SQL |\n", + "| `Trino` | `compute`: читает те же таблицы через общий каталог |\n", + "| `Jupyter` | Точка входа в учебные ноутбуки |\n", + "\n", + "В этом модуле мы не строим пайплайн и не создаём Iceberg-таблицы. Задача только одна: убедиться, что среда рабочая и понятная.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "module1-create-session", + "metadata": {}, + "outputs": [], + "source": [ + "from spark.cluster_smoke import (\n", + " create_spark_session,\n", + " format_smoke_report,\n", + " run_cluster_smoke,\n", + ")\n", + "\n", + "spark = create_spark_session(app_name=\"module-01-environment-and-smoke-test\")\n", + "spark\n" + ] + }, + { + "cell_type": "markdown", + "id": "module1-cluster-demo", + "metadata": {}, + "source": [ + "## Демонстрация: подключение к кластеру\n", + "\n", + "Сначала посмотрим на базовые признаки того, что ноутбук говорит именно с кластером, а не с локальным `local[*]` режимом.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "module1-connection-details", + "metadata": {}, + "outputs": [], + "source": [ + "connection_details = {\n", + " \"app_name\": spark.sparkContext.appName,\n", + " \"master_url\": spark.sparkContext.master,\n", + " \"default_parallelism\": spark.sparkContext.defaultParallelism,\n", + " \"spark_version\": spark.version,\n", + "}\n", + "connection_details\n" + ] + }, + { + "cell_type": "markdown", + "id": "module1-visual-check", + "metadata": {}, + "source": [ + "Проверь глазами:\n", + "\n", + "- `master_url` должен быть `spark://spark-master:7077`, а не `local[*]`;\n", + "- `default_parallelism` должен быть больше `1`;\n", + "- в `Spark UI` должны быть видны master и worker-ы.\n" + ] + }, + { + "cell_type": "markdown", + "id": "module1-smoke-explanation", + "metadata": {}, + "source": [ + "## Демонстрация: Spark smoke test\n", + "\n", + "Smoke test запускает распределённое вычисление на `range(0, 1_000_000)` и сверяет детерминированный результат.\n", + "\n", + "Если тест проходит, это означает минимум следующее:\n", + "\n", + "- Jupyter может создать SparkSession;\n", + "- Spark подключён к кластерному master;\n", + "- задание действительно выполняется, а не падает на первом действии.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "module1-run-smoke", + "metadata": {}, + "outputs": [], + "source": [ + "smoke_report = run_cluster_smoke(spark)\n", + "print(format_smoke_report(smoke_report))\n", + "smoke_report\n" + ] + }, + { + "cell_type": "markdown", + "id": "module1-self-check", + "metadata": {}, + "source": [ + "## Самостоятельное повторение\n", + "\n", + "Сделай руками и проверь себя:\n", + "\n", + "1. Открой `Spark UI`, `Trino UI`, `MinIO Console` и `JupyterLab`.\n", + "2. Сопоставь каждый сервис с одной из ролей: `storage`, `catalog`, `compute`, `entrypoint`.\n", + "3. Объясни, почему этот smoke test не должен работать в `local[*]` режиме.\n", + "4. На хосте выполни `docker compose ps` и посмотри, какие контейнеры должны быть в состоянии `Up`.\n", + "5. На хосте выполни `docker compose logs -f trino` или `docker compose logs -f spark-master`, чтобы посмотреть, где начинается первичная диагностика.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "module1-service-roles", + "metadata": {}, + "outputs": [], + "source": [ + "service_roles = {\n", + " # Заполни значения самостоятельно: storage, catalog, compute или entrypoint.\n", + " \"MinIO\": \"TODO\",\n", + " \"PostgreSQL\": \"TODO\",\n", + " \"Spark\": \"TODO\",\n", + " \"Trino\": \"TODO\",\n", + " \"Jupyter\": \"TODO\",\n", + "}\n", + "service_roles\n" + ] + }, + { + "cell_type": "markdown", + "id": "module1-checkpoint", + "metadata": {}, + "source": [ + "## Checkpoint\n", + "\n", + "К концу модуля ты должен уметь подтвердить:\n", + "\n", + "- стенд поднят и основные контейнеры живы;\n", + "- `Spark UI`, `MinIO Console`, `Trino UI` и `JupyterLab` открываются;\n", + "- smoke test завершился успешно и вернул ожидаемый результат;\n", + "- ты понимаешь назначение `MinIO`, `PostgreSQL`, `Spark`, `Trino`, `Jupyter`;\n", + "- ты знаешь, что первый маршрут диагностики начинается с `docker compose ps`, `docker compose logs -f ` и проверки UI.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "module1-stop-session", + "metadata": {}, + "outputs": [], + "source": [ + "spark.stop()\n" + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": "Python 3", + "language": "python", + "name": "python3" + }, + "language_info": { + "name": "python" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} diff --git a/notebooks/03_partitioning_and_schema_evolution.ipynb b/notebooks/03_partitioning_and_schema_evolution.ipynb deleted file mode 100755 index 52e5c05..0000000 --- a/notebooks/03_partitioning_and_schema_evolution.ipynb +++ /dev/null @@ -1,163 +0,0 @@ -{ - "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 deleted file mode 100755 index 623fb20..0000000 --- a/notebooks/spark-basic-test.ipynb +++ /dev/null @@ -1,86 +0,0 @@ -{ - "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/plans/README.md b/plans/README.md new file mode 100644 index 0000000..bcbac09 --- /dev/null +++ b/plans/README.md @@ -0,0 +1,16 @@ +# Планы работ + +Этот каталог хранит внутренние рабочие планы по развитию учебного курса. + +## Правила хранения + +- Формат: `living docs`. +- История изменений хранится в git. +- Именование файлов: `module-NN-.md`. +- В каждом плане фиксируются `Статус` и `Последнее обновление`. + +## Оглавление + +| Модуль | План | Статус | Целевые артефакты | Обновлён | +| --- | --- | --- | --- | --- | +| 1. Вход в стенд и базовая диагностика | [module-01-environment-and-basic-diagnostics.md](./module-01-environment-and-basic-diagnostics.md) | Ready for validation | `START_HERE.md`, `notebooks/01_environment_and_smoke_test.ipynb`, `src/spark/cluster_smoke.py` | 2026-03-06 | diff --git a/plans/module-01-environment-and-basic-diagnostics.md b/plans/module-01-environment-and-basic-diagnostics.md new file mode 100644 index 0000000..690dc1c --- /dev/null +++ b/plans/module-01-environment-and-basic-diagnostics.md @@ -0,0 +1,67 @@ +# Модуль 1. Вход в стенд и базовая диагностика + +**Статус:** `Ready for validation` +**Последнее обновление:** `2026-03-06` + +## Цель + +Дать студенту безопасный и воспроизводимый вход в курс: он должен сам поднять локальный стенд, открыть основные веб-интерфейсы, выполнить Spark smoke test и понять первый маршрут диагностики проблем. + +## Результат для студента + +После прохождения модуля студент: + +- поднимает стенд командами `docker compose build` и `docker compose up -d`; +- проверяет состояние контейнеров через `docker compose ps`; +- открывает `Spark UI`, `MinIO Console`, `Trino UI` и `JupyterLab`; +- запускает smoke test Spark-кластера и понимает, что именно он проверяет; +- знает, куда смотреть при первичной диагностике через логи и UI. + +## Deliverables + +- `plans/README.md` как индекс всех рабочих планов; +- `START_HERE.md` как канонический вход в курс до первого ноутбука; +- `notebooks/01_environment_and_smoke_test.ipynb` как канонический ноутбук Модуля 1; +- `src/spark/cluster_smoke.py` как общий helper и CLI smoke test; +- обновлённые ссылки и инструкции в `README.md` и `AGENTS.md`. + +## План работ + +1. Создать каталог `plans/` и зафиксировать в нём правила хранения и индекс. +2. Подготовить `START_HERE.md` с prerequisites, запуском, проверкой сервисов, входом в Jupyter и reset/restart шагами. +3. Расширить `src/spark/cluster_smoke.py` до reusable helper с функцией `run_cluster_smoke(spark) -> dict` и CLI `main()`. +4. Собрать `notebooks/01_environment_and_smoke_test.ipynb` по схеме `объяснение -> демонстрация -> самостоятельное повторение -> checkpoint`. +5. Перенести полезное содержимое из старого `spark-basic-test.ipynb` в новый ноутбук и убрать старый артефакт из активного маршрута. +6. Обновить `README.md` и `AGENTS.md`, чтобы новый вход и новый ноутбук стали каноническими. +7. Прогнать валидацию: smoke script должен проходить синтаксическую проверку, а ноутбук должен быть валидным `ipynb`. + +## Checkpoint + +Студент должен уметь: + +- показать результат `docker compose ps`; +- открыть `Spark UI`, `MinIO Console`, `Trino UI` и `JupyterLab`; +- выполнить smoke test и объяснить, что он подтвердил; +- назвать роли `MinIO`, `PostgreSQL`, `Spark`, `Trino` и `Jupyter`; +- описать первый шаг диагностики, если не поднимается `trino` или `spark-master`. + +## Acceptance Criteria + +- `START_HERE.md` позволяет поднять стенд с нуля без обращения к `docs/archive/legacy_howto.md`; +- `notebooks/01_environment_and_smoke_test.ipynb` можно выполнить сверху вниз в поднятом Jupyter-окружении; +- один и тот же smoke helper используется из CLI и из ноутбука; +- в репозитории остаётся один канонический вход в Модуль 1; +- `plans/README.md` позволяет быстро найти модульный план и понять его статус. + +## Риски + +- Старые ссылки на `spark-basic-test.ipynb` могут остаться в инструкциях и создавать путаницу. +- Если smoke test будет слишком сложным, он перестанет быть быстрым диагностическим шагом. +- Если в ноутбуке окажутся host-level команды, студент может ошибочно пытаться запускать `docker compose` внутри Jupyter. + +## Out of Scope + +- загрузка raw-данных; +- разбор internals `Iceberg`; +- запросы к `Trino` как отдельная практика; +- `schema evolution`, `time travel`, `compaction`, `vacuum`. diff --git a/src/spark/__init__.py b/src/spark/__init__.py new file mode 100644 index 0000000..2cb140b --- /dev/null +++ b/src/spark/__init__.py @@ -0,0 +1,3 @@ +from .cluster_smoke import create_spark_session, format_smoke_report, run_cluster_smoke + +__all__ = ["create_spark_session", "format_smoke_report", "run_cluster_smoke"] diff --git a/src/spark/cluster_smoke.py b/src/spark/cluster_smoke.py index a222785..2a81a73 100755 --- a/src/spark/cluster_smoke.py +++ b/src/spark/cluster_smoke.py @@ -1,12 +1,84 @@ +from __future__ import annotations + from pyspark.sql import SparkSession -spark = ( - SparkSession.builder - .appName("test-cluster") - .master("spark://spark-master:7077") # важно: не local[*] - .getOrCreate() -) +DEFAULT_MASTER_URL = "spark://spark-master:7077" +DEFAULT_APP_NAME = "cluster-smoke" +DEFAULT_ROW_COUNT = 1_000_000 -spark.range(0, 1000000).groupBy().sum().show() -spark.stop() +def create_spark_session( + app_name: str = DEFAULT_APP_NAME, + master_url: str = DEFAULT_MASTER_URL, +) -> SparkSession: + """Создать SparkSession для проверок окружения в Модуле 1. + + Функция специально остаётся минимальной и опирается на смонтированный + в контейнер `spark-defaults.conf` для остальной Spark/Iceberg-конфигурации. + """ + return SparkSession.builder.appName(app_name).master(master_url).getOrCreate() + + +def run_cluster_smoke(spark: SparkSession, row_count: int = DEFAULT_ROW_COUNT) -> dict: + spark_context = spark.sparkContext + master_url = spark_context.master + + if master_url.startswith("local"): + raise RuntimeError( + "Smoke test must run against the Spark cluster, not local mode." + ) + + # Берём больше партиций, чем даёт defaultParallelism, чтобы задача точно + # распределилась по worker-ам, а не схлопнулась в слишком маленький job. + partition_count = max(spark_context.defaultParallelism * 2, 8) + expected_sum = row_count * (row_count - 1) // 2 + actual_sum = ( + spark.range(0, row_count, 1, numPartitions=partition_count) + .selectExpr("SUM(id) AS total_sum") + .collect()[0]["total_sum"] + ) + + if actual_sum != expected_sum: + raise RuntimeError( + f"Unexpected smoke result: expected {expected_sum}, got {actual_sum}." + ) + + return { + "app_name": spark_context.appName, + "master_url": master_url, + "spark_version": spark.version, + "default_parallelism": spark_context.defaultParallelism, + "partition_count": partition_count, + "row_count": row_count, + "result_sum": actual_sum, + "expected_sum": expected_sum, + } + + +def format_smoke_report(report: dict) -> str: + return "\n".join( + [ + "Spark cluster smoke test passed.", + f"app_name={report['app_name']}", + f"master_url={report['master_url']}", + f"spark_version={report['spark_version']}", + f"default_parallelism={report['default_parallelism']}", + f"partition_count={report['partition_count']}", + f"row_count={report['row_count']}", + f"result_sum={report['result_sum']}", + f"expected_sum={report['expected_sum']}", + ] + ) + + +def main() -> None: + spark = create_spark_session() + try: + report = run_cluster_smoke(spark) + print(format_smoke_report(report)) + finally: + spark.stop() + + +if __name__ == "__main__": + main()