diff --git a/.gitignore b/.gitignore index b3ea1eb..88e0a28 100755 --- a/.gitignore +++ b/.gitignore @@ -21,3 +21,9 @@ venv/ # Logs *.log + +# Dataset (downloaded by student, not stored in repo) +/data/* +!/data/nyc_taxi/ +/data/nyc_taxi/* +!/data/nyc_taxi/.gitkeep diff --git a/START_HERE.md b/START_HERE.md index efebd2d..ddb61d6 100644 --- a/START_HERE.md +++ b/START_HERE.md @@ -1,6 +1,6 @@ # START HERE -Этот документ нужен для первого входа в курс и прохождения Модуля 1: поднять стенд, проверить сервисы, открыть интерфейсы и выполнить базовый smoke test. +Этот документ нужен для первого входа в курс: поднять стенд, проверить сервисы, открыть интерфейсы, выполнить базовый smoke test и подготовить локальный data bundle для следующих модулей. ## Маршрут прохождения @@ -9,7 +9,8 @@ 3. Убедиться, что контейнеры живы. 4. Открыть основные UI. 5. Зайти в Jupyter и выполнить `notebooks/01_environment_and_smoke_test.ipynb`. -6. При проблемах использовать логи и шаги диагностики из этого документа. +6. Перед Модулем 3 подготовить учебный датасет в `./data/nyc_taxi`. +7. При проблемах использовать логи и шаги диагностики из этого документа. ## Prerequisites @@ -100,8 +101,10 @@ notebooks/01_environment_and_smoke_test.ipynb - `./notebooks` смонтирован в контейнер как `/opt/work`; - `./src` смонтирован read-only как `/opt/src`; +- `./data` смонтирован read-only как `/opt/data`; - `PYTHONPATH=/opt/src`, поэтому helper-скрипты из `src/` доступны для импорта в ноутбуках; -- первый ноутбук курса: `01_environment_and_smoke_test.ipynb`. +- первый ноутбук курса: `01_environment_and_smoke_test.ipynb`; +- локальный учебный data bundle для Модуля 3 должен быть доступен внутри Jupyter по пути `/opt/data/nyc_taxi`. ## Базовая диагностика @@ -146,8 +149,42 @@ docker compose up -d Используй полный reset, если хочешь пройти практику заново с чистого состояния. +## Подготовка учебного датасета (перед Модулем 3) + +Перед `notebooks/03_raw_ingest_and_first_read.ipynb` нужно заранее скачать учебный data bundle на хост, а не изнутри ноутбука. + +Скачай каноническое учебное подмножество `NYC TLC Yellow Taxi Trip Records`: + +```bash +curl -fLo data/nyc_taxi/yellow_tripdata_2024-01.parquet \ + https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2024-01.parquet +curl -fLo data/nyc_taxi/yellow_tripdata_2024-02.parquet \ + https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2024-02.parquet +curl -fLo data/nyc_taxi/yellow_tripdata_2024-03.parquet \ + https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2024-03.parquet +curl -fLo data/nyc_taxi/taxi_zone_lookup.csv \ + https://d37ci6vzurychx.cloudfront.net/misc/taxi_zone_lookup.csv +``` + +Проверь, что файлы на месте: + +```bash +ls -lh data/nyc_taxi/ +``` + +Минимальный набор для курса: + +- `yellow_tripdata_2024-01.parquet` +- `yellow_tripdata_2024-02.parquet` +- `yellow_tripdata_2024-03.parquet` +- `taxi_zone_lookup.csv` + +Если хочешь расширенный режим, можешь скачать все 12 месяцев `2024`, но основной маршрут курса и примеры опираются на первые 3 месяца. + ## Что делать дальше - пройти `notebooks/01_environment_and_smoke_test.ipynb`; +- пройти `notebooks/02_lakehouse_mental_model.ipynb`; +- подготовить data bundle по инструкции выше и пройти `notebooks/03_raw_ingest_and_first_read.ipynb`; - свериться с программой курса в `docs/course_program.md`; -- после прохождения Модуля 1 переходить к следующим учебным материалам. +- после прохождения первых модулей переходить к следующим учебным материалам. diff --git a/data/nyc_taxi/.gitkeep b/data/nyc_taxi/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/docker-compose.yml b/docker-compose.yml index 39451c7..48596ab 100755 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -181,6 +181,7 @@ services: - ./spark/spark-defaults.conf:/opt/spark/conf/spark-defaults.conf:ro - ./notebooks:/opt/work - ./src:/opt/src:ro + - ./data:/opt/data:ro trino: image: trinodb/trino:478 diff --git a/docs/stack_reference.md b/docs/stack_reference.md index 1c078b6..b5c5154 100644 --- a/docs/stack_reference.md +++ b/docs/stack_reference.md @@ -110,10 +110,13 @@ docker compose run --rm iceberg-catalog-init - адрес: `http://localhost:8888` - `./notebooks` смонтирован как `/opt/work` - `./src` смонтирован read-only как `/opt/src` +- `./data` смонтирован read-only как `/opt/data` - `PYTHONPATH=/opt/src` Первый ноутбук курса: `notebooks/01_environment_and_smoke_test.ipynb`. +Для Модуля 3 локально скачанный data bundle должен лежать на хосте в `./data/nyc_taxi` и будет доступен внутри Jupyter по пути `/opt/data/nyc_taxi`. + ## Как связаны Spark, Trino, PostgreSQL и MinIO - `Spark` использует каталог `lakehouse` через `JdbcCatalog`. diff --git a/notebooks/03_raw_ingest_and_first_read.ipynb b/notebooks/03_raw_ingest_and_first_read.ipynb new file mode 100644 index 0000000..e867b9c --- /dev/null +++ b/notebooks/03_raw_ingest_and_first_read.ipynb @@ -0,0 +1,685 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "id": "471a32f6", + "metadata": {}, + "source": [ + "# Модуль 3. Raw-данные и первое чтение в Spark\n", + "\n", + "В этом модуле мы разбираем полный первый цикл работы с учебным датасетом в Lakehouse:\n", + "\n", + "- проверяем заранее скачанный локальный data bundle;\n", + "- загружаем исходные файлы в raw-зону `MinIO`;\n", + "- читаем raw-данные через `Spark` по `s3a://` пути;\n", + "- смотрим схему, типы, базовые метрики и простые аномалии;\n", + "- фиксируем, почему raw-слой нужен как воспроизводимая неизменяемая точка входа.\n", + "\n", + "После этого ноутбука ты должен уметь объяснить, почему raw-данные не надо «чинить на месте» и почему их не стоит скачивать прямо из ноутбука.\n" + ] + }, + { + "cell_type": "markdown", + "id": "723dba8f", + "metadata": {}, + "source": [ + "## 0. Перед стартом\n", + "\n", + "Перед выполнением ноутбука:\n", + "\n", + "- подними стенд по `START_HERE.md`;\n", + "- пройди `notebooks/01_environment_and_smoke_test.ipynb`;\n", + "- желательно пройти `notebooks/02_lakehouse_mental_model.ipynb`;\n", + "- заранее скачай data bundle в `./data/nyc_taxi` по инструкции из `START_HERE.md`.\n", + "\n", + "В классическом `PostgreSQL / Greenplum` данные часто уже лежат внутри системы. В Lakehouse путь более явный: сначала исходные файлы должны попасть в object storage, и только потом мы начинаем поверх них читать, профилировать и строить таблицы.\n", + "\n", + "| | Классический DWH | Lakehouse |\n", + "| --- | --- | --- |\n", + "| Где стартует работа | Данные уже внутри СУБД | Исходные файлы сначала попадают в storage |\n", + "| Первая практическая задача | `SELECT` из готовой таблицы | Загрузить raw и прочитать raw |\n", + "| Что важно не сломать | Таблицу / схему БД | Неизменяемый raw-слой |\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "a5c56637", + "metadata": {}, + "outputs": [], + "source": [ + "from pathlib import Path\n", + "import re\n", + "from datetime import datetime\n", + "\n", + "import boto3\n", + "import pandas as pd\n", + "from botocore.exceptions import ClientError\n", + "from pyspark.sql import SparkSession, functions as F\n", + "\n", + "DATA_DIR = Path(\"/opt/data/nyc_taxi\")\n", + "BUCKET_NAME = \"lakehouse\"\n", + "RAW_PREFIX = \"raw/nyc_taxi/\"\n", + "RAW_URI = f\"s3a://{BUCKET_NAME}/{RAW_PREFIX}\"\n", + "ZONE_LOOKUP_RAW_URI = f\"{RAW_URI}taxi_zone_lookup.csv\"\n", + "\n", + "spark = SparkSession.builder .appName(\"module-03-raw-ingest-and-first-read\") .getOrCreate()\n", + "\n", + "spark.sparkContext.setLogLevel(\"ERROR\")\n", + "\n", + "s3 = boto3.client(\n", + " \"s3\",\n", + " endpoint_url=\"http://minio:9000\",\n", + " aws_access_key_id=\"minioadmin\",\n", + " aws_secret_access_key=\"minioadmin\",\n", + ")\n", + "\n", + "\n", + "def format_bytes(num_bytes: int) -> str:\n", + " units = [\"B\", \"KB\", \"MB\", \"GB\", \"TB\"]\n", + " value = float(num_bytes)\n", + " for unit in units:\n", + " if value < 1024 or unit == units[-1]:\n", + " return f\"{value:.1f} {unit}\"\n", + " value /= 1024\n", + "\n", + "\n", + "def object_exists(bucket: str, key: str) -> bool:\n", + " try:\n", + " s3.head_object(Bucket=bucket, Key=key)\n", + " return True\n", + " except ClientError as exc:\n", + " error_code = exc.response.get(\"Error\", {}).get(\"Code\", \"\")\n", + " if error_code in {\"404\", \"NoSuchKey\", \"NotFound\"}:\n", + " return False\n", + " raise\n", + "\n", + "\n", + "def list_objects(bucket: str, prefix: str) -> list[dict]:\n", + " paginator = s3.get_paginator(\"list_objects_v2\")\n", + " objects = []\n", + " for page in paginator.paginate(Bucket=bucket, Prefix=prefix):\n", + " objects.extend(page.get(\"Contents\", []))\n", + " return objects\n", + "\n", + "\n", + "def infer_month_bounds(paths: list[Path]) -> tuple[datetime, datetime]:\n", + " months = []\n", + " for path in paths:\n", + " match = re.fullmatch(r\"yellow_tripdata_(\\d{4})-(\\d{2})\\.parquet\", path.name)\n", + " if match:\n", + " year = int(match.group(1))\n", + " month = int(match.group(2))\n", + " months.append((year, month))\n", + "\n", + " assert months, \"Не удалось определить месяцы из имён yellow_tripdata_*.parquet\"\n", + "\n", + " start_year, start_month = min(months)\n", + " end_year, end_month = max(months)\n", + "\n", + " expected_start = datetime(start_year, start_month, 1)\n", + " if end_month == 12:\n", + " expected_end = datetime(end_year + 1, 1, 1)\n", + " else:\n", + " expected_end = datetime(end_year, end_month + 1, 1)\n", + "\n", + " return expected_start, expected_end\n", + "\n", + "\n", + "print(\"Spark и MinIO-клиент готовы\")\n", + "print(f\"RAW URI: {RAW_URI}\")\n" + ] + }, + { + "cell_type": "markdown", + "id": "d4d8f302", + "metadata": {}, + "source": [ + "## 1. Проверка data bundle\n", + "\n", + "Сначала убеждаемся, что нужные файлы уже лежат внутри Jupyter-контейнера по пути `/opt/data/nyc_taxi`.\n", + "\n", + "Почему не скачиваем данные прямо здесь:\n", + "\n", + "- так модуль остаётся воспроизводимым и не зависит от сетевого доступа внутри ноутбука;\n", + "- onboarding и подготовка окружения остаются в `START_HERE.md`, а практика в ноутбуке занимается только Lakehouse-операциями;\n", + "- студент видит границу ответственности: сначала подготовить bundle, потом работать с raw-слоем.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "da2f7b51", + "metadata": {}, + "outputs": [], + "source": [ + "assert DATA_DIR.exists(), (\n", + " \"Каталог /opt/data/nyc_taxi не найден. \"\n", + " \"Сначала скачай data bundle по инструкции из START_HERE.md и убедись, что ./data смонтирован в Jupyter как /opt/data.\"\n", + ")\n", + "\n", + "bundle_files = sorted(path for path in DATA_DIR.iterdir() if path.is_file())\n", + "yellow_files = sorted(DATA_DIR.glob(\"yellow_tripdata_*.parquet\"))\n", + "zone_lookup_file = DATA_DIR / \"taxi_zone_lookup.csv\"\n", + "\n", + "assert yellow_files, (\n", + " \"Не найдены файлы yellow_tripdata_*.parquet в /opt/data/nyc_taxi. \"\n", + " \"Скачай каноническое подмножество из START_HERE.md.\"\n", + ")\n", + "assert zone_lookup_file.exists(), (\n", + " \"Не найден taxi_zone_lookup.csv. Он нужен и для этого модуля, и для самостоятельного задания. \"\n", + " \"Скачай его по инструкции из START_HERE.md.\"\n", + ")\n", + "\n", + "bundle_report = pd.DataFrame(\n", + " [\n", + " {\n", + " \"file\": path.name,\n", + " \"size_bytes\": path.stat().st_size,\n", + " \"size_human\": format_bytes(path.stat().st_size),\n", + " }\n", + " for path in bundle_files\n", + " ]\n", + ").sort_values(\"file\").reset_index(drop=True)\n", + "\n", + "bundle_report\n" + ] + }, + { + "cell_type": "markdown", + "id": "c8defd34", + "metadata": {}, + "source": [ + "Набор файлов должен быть понятным и фиксированным. В этом модуле мы работаем не с «каким-то случайным parquet из интернета», а с заранее определённым учебным набором.\n" + ] + }, + { + "cell_type": "markdown", + "id": "dabd90f2", + "metadata": {}, + "source": [ + "## 2. Загрузка в raw-зону MinIO\n", + "\n", + "Raw-слой в Lakehouse похож на воспроизводимый входной `stg`-контур: мы сохраняем исходные файлы как есть, под оригинальными именами, не переписывая их содержимое.\n", + "\n", + "Конвенция пути в этом курсе:\n", + "\n", + "- bucket: `lakehouse`\n", + "- raw prefix: `raw/nyc_taxi/`\n", + "- полный путь для чтения: `s3a://lakehouse/raw/nyc_taxi/`\n", + "\n", + "Важно: raw-зона не смешивается с `warehouse/`, где потом будут жить Iceberg-таблицы.\n", + "\n", + "Повторный запуск этой секции должен быть безопасным:\n", + "\n", + "- если объекта ещё нет в raw-зоне, мы его загружаем;\n", + "- если объект уже есть и размер совпадает, мы пропускаем загрузку;\n", + "- если объект уже есть, но размер отличается, это сигнал проблемы с данными или bundle, и мы падаем сразу, без перезаписи raw.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "8644cc94", + "metadata": {}, + "outputs": [], + "source": [ + "upload_candidates = [*yellow_files, zone_lookup_file]\n", + "upload_report = []\n", + "\n", + "for path in upload_candidates:\n", + " object_key = f\"{RAW_PREFIX}{path.name}\"\n", + " local_size = path.stat().st_size\n", + " existed_before = object_exists(BUCKET_NAME, object_key)\n", + "\n", + " if existed_before:\n", + " metadata = s3.head_object(Bucket=BUCKET_NAME, Key=object_key)\n", + " remote_size = metadata[\"ContentLength\"]\n", + " if remote_size != local_size:\n", + " raise ValueError(\n", + " \"Raw object already exists with different size: \"\n", + " f\"{object_key} (local={local_size}, remote={remote_size}). \"\n", + " \"Остановись и проверь data bundle вместо перезаписи raw-слоя.\"\n", + " )\n", + " action = \"skipped_same_size\"\n", + " else:\n", + " s3.upload_file(str(path), BUCKET_NAME, object_key)\n", + " metadata = s3.head_object(Bucket=BUCKET_NAME, Key=object_key)\n", + " remote_size = metadata[\"ContentLength\"]\n", + " action = \"uploaded\"\n", + "\n", + " upload_report.append(\n", + " {\n", + " \"file\": path.name,\n", + " \"raw_key\": object_key,\n", + " \"action\": action,\n", + " \"local_size_human\": format_bytes(local_size),\n", + " \"remote_size_human\": format_bytes(remote_size),\n", + " }\n", + " )\n", + "\n", + "pd.DataFrame(upload_report)\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "11f4c459", + "metadata": {}, + "outputs": [], + "source": [ + "raw_objects = list_objects(BUCKET_NAME, RAW_PREFIX)\n", + "\n", + "pd.DataFrame(\n", + " [\n", + " {\n", + " \"key\": obj[\"Key\"],\n", + " \"size_human\": format_bytes(obj[\"Size\"]),\n", + " }\n", + " for obj in raw_objects\n", + " ]\n", + ")\n" + ] + }, + { + "cell_type": "markdown", + "id": "48295c73", + "metadata": {}, + "source": [ + "Если объект уже существовал, мы не перезаписываем его без необходимости. Совпадающий размер означает, что повторный запуск секции идемпотентен. Несовпадающий размер считается конфликтом и останавливает сценарий до любых изменений в raw.\n", + "\n", + "На этом шаге данные уже лежат в storage. `Spark` может читать их напрямую по `s3a://` пути, а для `Trino` понадобится зарегистрированная таблица. До неё мы дойдём в Модулях 4 и 6.\n" + ] + }, + { + "cell_type": "markdown", + "id": "696d6adc", + "metadata": {}, + "source": [ + "## 3. Первое чтение raw PARQUET через Spark\n", + "\n", + "Теперь читаем raw-файлы напрямую из `MinIO` по `s3a://` пути. Это и есть первая важная практика Lakehouse: вычислитель (`Spark`) работает поверх файлов в object storage, а не только поверх локальной файловой системы или таблиц внутри одной СУБД.\n", + "\n", + "Здесь важно различать два режима работы:\n", + "\n", + "- `spark.read.parquet(\"s3a://...\")` читает raw-файлы напрямую по физическому пути;\n", + "- `spark.table(\"lakehouse....\")` читает уже зарегистрированную управляемую таблицу через каталог.\n", + "\n", + "Сейчас мы ещё не создавали `Iceberg`-таблицу, поэтому работаем именно через `spark.read.parquet(...)`. К `spark.table(...)` вернёмся в Модуле 4.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "0e8823aa", + "metadata": {}, + "outputs": [], + "source": [ + "single_file_uri = f\"{RAW_URI}yellow_tripdata_2024-01.parquet\"\n", + "single_month_df = spark.read.parquet(single_file_uri)\n", + "\n", + "print(f\"Читаем один raw-файл: {single_file_uri}\")\n", + "print(f\"Строк в одном файле: {single_month_df.count():,}\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "fd5ad623", + "metadata": {}, + "outputs": [], + "source": [ + "yellow_df = spark.read.parquet(f\"{RAW_URI}yellow_tripdata_*.parquet\")\n", + "row_count = yellow_df.count()\n", + "\n", + "print(f\"Прочитано строк: {row_count:,}\")\n", + "print(f\"Количество parquet-файлов в наборе: {len(yellow_files)}\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "5e147271", + "metadata": {}, + "outputs": [], + "source": [ + "yellow_df.show(10, truncate=False)\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "efa037e6", + "metadata": {}, + "outputs": [], + "source": [ + "input_files = sorted(yellow_df.inputFiles())\n", + "pd.DataFrame({\"input_file\": input_files})\n" + ] + }, + { + "cell_type": "markdown", + "id": "e8866f38", + "metadata": {}, + "source": [ + "Обрати внимание: мы читаем сразу несколько файлов wildcard-путём. Это обычный рабочий сценарий для raw-слоя, где данные часто приходят партиями по датам, дням или месяцам.\n", + "\n", + "В этом модуле мы намеренно не используем `cache()` для всего raw DataFrame. Для учебного профилирования важнее устойчиво прочитать исходные файлы, чем пытаться держать весь raw-срез в памяти executor-ов.\n" + ] + }, + { + "cell_type": "markdown", + "id": "9701728a", + "metadata": {}, + "source": [ + "## 4. Схема, типы и null-значения\n", + "\n", + "Следующий шаг после первого чтения: понять, какие колонки пришли, какие у них типы и где уже на raw-слое встречаются `null`.\n", + "\n", + "Здесь важно не перепутать диагностику и очистку. В raw-слое мы наблюдаем и фиксируем реальность данных, а не исправляем её.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "e2809a8b", + "metadata": {}, + "outputs": [], + "source": [ + "yellow_df.printSchema()\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "126e5601", + "metadata": {}, + "outputs": [], + "source": [ + "schema_df = pd.DataFrame(yellow_df.dtypes, columns=[\"column\", \"spark_type\"])\n", + "print(f\"Количество колонок: {len(yellow_df.columns)}\")\n", + "schema_df\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "c6093e2a", + "metadata": {}, + "outputs": [], + "source": [ + "null_counts_row = yellow_df.select(\n", + " *[F.sum(F.col(column).isNull().cast(\"int\")).alias(column) for column in yellow_df.columns]\n", + ").toPandas().iloc[0]\n", + "\n", + "null_report = pd.DataFrame(\n", + " {\n", + " \"column\": yellow_df.columns,\n", + " \"null_count\": [int(null_counts_row[column]) for column in yellow_df.columns],\n", + " }\n", + ")\n", + "null_report[\"null_share\"] = (null_report[\"null_count\"] / row_count).round(6)\n", + "null_report.sort_values([\"null_count\", \"column\"], ascending=[False, True]).reset_index(drop=True)\n" + ] + }, + { + "cell_type": "markdown", + "id": "567156e6", + "metadata": {}, + "source": [ + "`Null` в raw-данных сам по себе не является багом загрузки. Это нормальная часть профиля источника. Позже, на `bronze` и `silver`, мы будем принимать явные решения: какие поля очищать, что отбрасывать, а что нормализовать.\n" + ] + }, + { + "cell_type": "markdown", + "id": "69e6b77e", + "metadata": {}, + "source": [ + "## 5. Базовые метрики и аномалии\n", + "\n", + "Теперь быстро профилируем данные и ищем несколько очевидных проблемных паттернов:\n", + "\n", + "- отрицательный `fare_amount`;\n", + "- нулевая дистанция при ненулевой сумме;\n", + "- слишком большой `total_amount`;\n", + "- `pickup` вне ожидаемого диапазона дат.\n", + "\n", + "Диапазон дат будем выводить из имён файлов, а не хардкодить руками.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "3a4ed596", + "metadata": {}, + "outputs": [], + "source": [ + "numeric_columns = [\n", + " column\n", + " for column, spark_type in yellow_df.dtypes\n", + " if spark_type.startswith((\"tinyint\", \"smallint\", \"int\", \"bigint\", \"float\", \"double\", \"decimal\"))\n", + "]\n", + "\n", + "yellow_df.select(*numeric_columns).describe().toPandas()\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "cc029c98", + "metadata": {}, + "outputs": [], + "source": [ + "required_columns = {\n", + " \"fare_amount\",\n", + " \"trip_distance\",\n", + " \"total_amount\",\n", + " \"tpep_pickup_datetime\",\n", + " \"PULocationID\",\n", + " \"DOLocationID\",\n", + "}\n", + "missing_columns = sorted(required_columns.difference(yellow_df.columns))\n", + "assert not missing_columns, (\n", + " \"В raw-схеме не хватает ожидаемых колонок: \"\n", + " f\"{missing_columns}. Возможно, источник поменял формат, и модуль требует обновления.\"\n", + ")\n", + "\n", + "expected_start, expected_end = infer_month_bounds(yellow_files)\n", + "\n", + "anomaly_report = pd.DataFrame(\n", + " [\n", + " {\n", + " \"check\": \"fare_amount < 0\",\n", + " \"count\": yellow_df.filter(F.col(\"fare_amount\") < 0).count(),\n", + " },\n", + " {\n", + " \"check\": \"trip_distance = 0 and total_amount != 0\",\n", + " \"count\": yellow_df.filter(\n", + " (F.col(\"trip_distance\") == 0) & (F.col(\"total_amount\") != 0)\n", + " ).count(),\n", + " },\n", + " {\n", + " \"check\": \"total_amount > 1000\",\n", + " \"count\": yellow_df.filter(F.col(\"total_amount\") > 1000).count(),\n", + " },\n", + " {\n", + " \"check\": f\"pickup outside [{expected_start.date()}, {expected_end.date()})\",\n", + " \"count\": yellow_df.filter(\n", + " (F.col(\"tpep_pickup_datetime\") < F.lit(expected_start))\n", + " | (F.col(\"tpep_pickup_datetime\") >= F.lit(expected_end))\n", + " ).count(),\n", + " },\n", + " ]\n", + ")\n", + "\n", + "anomaly_report\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "0cc5aaf7", + "metadata": {}, + "outputs": [], + "source": [ + "suspicious_rows = yellow_df.filter(\n", + " (F.col(\"fare_amount\") < 0)\n", + " | ((F.col(\"trip_distance\") == 0) & (F.col(\"total_amount\") != 0))\n", + " | (F.col(\"total_amount\") > 1000)\n", + " | (F.col(\"tpep_pickup_datetime\") < F.lit(expected_start))\n", + " | (F.col(\"tpep_pickup_datetime\") >= F.lit(expected_end))\n", + ")\n", + "\n", + "suspicious_rows.select(\n", + " \"tpep_pickup_datetime\",\n", + " \"passenger_count\",\n", + " \"trip_distance\",\n", + " \"fare_amount\",\n", + " \"total_amount\",\n", + " \"PULocationID\",\n", + " \"DOLocationID\",\n", + ").show(20, truncate=False)\n" + ] + }, + { + "cell_type": "markdown", + "id": "9ac27204", + "metadata": {}, + "source": [ + "Реальные данные почти всегда содержат выбросы, странные записи и неоднозначные значения. Это не повод немедленно переписывать raw-файлы. Наоборот, raw нужен именно затем, чтобы исходный срез можно было перечитать и переработать заново с новыми правилами.\n" + ] + }, + { + "cell_type": "markdown", + "id": "095c9ac6", + "metadata": {}, + "source": [ + "## 6. Самостоятельное задание\n", + "\n", + "Используй raw-данные, которые уже лежат в `MinIO`:\n", + "\n", + "1. Прочитай `taxi_zone_lookup.csv` из raw-зоны через `Spark`.\n", + "2. Выведи схему и первые 10 строк.\n", + "3. Посчитай количество зон по `Borough`.\n", + "4. Найди `top-5` `pickup`-локаций в поездках и присоедини к ним названия зон.\n", + "5. Найди ещё одну аномалию или интересный паттерн и коротко опиши его.\n", + "\n", + "Подсказка: CSV-файл лежит по пути `s3a://lakehouse/raw/nyc_taxi/taxi_zone_lookup.csv`.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "fe530d7d", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: прочитай taxi_zone_lookup.csv из raw-зоны через Spark\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "99fb5256", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: выведи схему и первые 10 строк\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "219a40fc", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: посчитай количество зон по Borough и top-5 pickup locations с join на lookup\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "5c944f60", + "metadata": {}, + "outputs": [], + "source": [ + "# Ваш код: найди ещё одну аномалию или интересный паттерн в данных\n" + ] + }, + { + "cell_type": "markdown", + "id": "e23dfe47", + "metadata": {}, + "source": [ + "## 7. Почему raw-слой важен\n", + "\n", + "Raw-слой полезен не потому, что в нём удобно жить постоянно, а потому что он даёт устойчивую исходную точку. Если завтра поменяются правила очистки, дедупликации или нормализации, мы не обязаны заново скачивать источник и гадать, что именно уже успели «подправить». Мы просто перечитываем raw и перестраиваем downstream-слои.\n", + "\n", + "Для инженера с опытом `PostgreSQL / Greenplum` полезна такая аналогия: raw похож на аккуратный внешний `staging`-контур, но в Lakehouse он обычно хранится в object storage как набор файлов, а не как таблица внутри одной СУБД.\n", + "\n", + "Отдельно важно, что загрузка источника вынесена из ноутбука в onboarding. Это делает практику менее хрупкой: ноутбук не зависит от внешней сети, а набор данных остаётся одинаковым у всех студентов.\n", + "\n", + "Наконец, raw не стоит чинить «на месте». Если в исходных данных есть `null`, отрицательные значения или странные даты, это сигнал для явных трансформаций в `bronze` и `silver`, а не повод тихо переписать входные файлы.\n" + ] + }, + { + "cell_type": "markdown", + "id": "8849c73f", + "metadata": {}, + "source": [ + "## 8. Checkpoint\n", + "\n", + "Проверь себя:\n", + "\n", + "1. Можешь ли ты показать файлы raw-зоны в `MinIO` через код или `MinIO Console`?\n", + "2. Сколько колонок в схеме `Yellow Taxi` и какие типы там встречаются?\n", + "3. Какие колонки содержат `null` и почему это не обязательно баг загрузки?\n", + "4. Зачем нужен raw-слой, если дальше всё равно будут трансформации?\n", + "5. Почему мы не скачиваем датасет прямо из ноутбука?\n", + "6. Что произойдёт, если источник завтра поменяет формат файлов?\n" + ] + }, + { + "cell_type": "markdown", + "id": "371c4c9d", + "metadata": {}, + "source": [ + "## 9. Завершение\n", + "\n", + "Мы намеренно не удаляем raw-данные из `MinIO`. В Модуле 4 они понадобятся для построения первого управляемого `bronze`-слоя на `Iceberg`.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "8574e96b", + "metadata": {}, + "outputs": [], + "source": [ + "spark.stop()\n" + ] + } + ], + "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 index a843ae2..41e26b5 100644 --- a/plans/README.md +++ b/plans/README.md @@ -15,3 +15,4 @@ | --- | --- | --- | --- | --- | | 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 | | 2. Ментальная модель Lakehouse: storage, catalog, compute | [module-02-lakehouse-mental-model.md](./module-02-lakehouse-mental-model.md) | Ready for validation | `notebooks/02_lakehouse_mental_model.ipynb` | 2026-03-07 | +| 3. Raw-данные и первое чтение в Spark | [module-03-raw-ingest-and-first-read.md](./module-03-raw-ingest-and-first-read.md) | Ready for validation | `docker-compose.yml`, `.gitignore`, `START_HERE.md`, `docs/stack_reference.md`, `notebooks/03_raw_ingest_and_first_read.ipynb` | 2026-03-07 | diff --git a/plans/module-03-raw-ingest-and-first-read.md b/plans/module-03-raw-ingest-and-first-read.md index ad317f2..1442f8c 100644 --- a/plans/module-03-raw-ingest-and-first-read.md +++ b/plans/module-03-raw-ingest-and-first-read.md @@ -1,6 +1,6 @@ # Модуль 3. Raw-данные и первое чтение в Spark -**Статус:** `Draft` +**Статус:** `Ready for validation` **Последнее обновление:** `2026-03-07` ## Цель