feat(course): добавлен модуль 3 про raw ingest и чтение в Spark

- Зачем:
  - нужен следующий практический шаг после модулей 1-2: доставить raw-данные в MinIO и впервые прочитать их через Spark.
- Что:
  - добавлен ноутбук модуля 3 с raw ingest, first read, проверкой схемы, null-анализом и поиском аномалий.
  - обновлены START_HERE, stack reference и docker-compose для data bundle в ./data и mount в /opt/data.
  - добавлены правила игнорирования data bundle, .gitkeep для пустой директории и обновлён статус плана модуля.
- Проверка:
  - docker compose config.
  - выполнение notebooks/03_raw_ingest_and_first_read.ipynb на локальном data bundle.
This commit is contained in:
2026-03-07 18:51:45 +03:00
parent 7ae771df68
commit 6d1b6e9f61
8 changed files with 738 additions and 5 deletions
+6
View File
@@ -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
+41 -4
View File
@@ -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 переходить к следующим учебным материалам.
- после прохождения первых модулей переходить к следующим учебным материалам.
View File
+1
View File
@@ -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
+3
View File
@@ -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`.
@@ -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
}
+1
View File
@@ -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 |
+1 -1
View File
@@ -1,6 +1,6 @@
# Модуль 3. Raw-данные и первое чтение в Spark
**Статус:** `Draft`
**Статус:** `Ready for validation`
**Последнее обновление:** `2026-03-07`
## Цель