{ "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 }