- Зачем: - нужен следующий практический шаг после модулей 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.
29 KiB
29 KiB
In [ ]:
from pathlib import Path
import re
from datetime import datetime
import boto3
import pandas as pd
from botocore.exceptions import ClientError
from pyspark.sql import SparkSession, functions as F
DATA_DIR = Path("/opt/data/nyc_taxi")
BUCKET_NAME = "lakehouse"
RAW_PREFIX = "raw/nyc_taxi/"
RAW_URI = f"s3a://{BUCKET_NAME}/{RAW_PREFIX}"
ZONE_LOOKUP_RAW_URI = f"{RAW_URI}taxi_zone_lookup.csv"
spark = SparkSession.builder .appName("module-03-raw-ingest-and-first-read") .getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
s3 = boto3.client(
"s3",
endpoint_url="http://minio:9000",
aws_access_key_id="minioadmin",
aws_secret_access_key="minioadmin",
)
def format_bytes(num_bytes: int) -> str:
units = ["B", "KB", "MB", "GB", "TB"]
value = float(num_bytes)
for unit in units:
if value < 1024 or unit == units[-1]:
return f"{value:.1f} {unit}"
value /= 1024
def object_exists(bucket: str, key: str) -> bool:
try:
s3.head_object(Bucket=bucket, Key=key)
return True
except ClientError as exc:
error_code = exc.response.get("Error", {}).get("Code", "")
if error_code in {"404", "NoSuchKey", "NotFound"}:
return False
raise
def list_objects(bucket: str, prefix: str) -> list[dict]:
paginator = s3.get_paginator("list_objects_v2")
objects = []
for page in paginator.paginate(Bucket=bucket, Prefix=prefix):
objects.extend(page.get("Contents", []))
return objects
def infer_month_bounds(paths: list[Path]) -> tuple[datetime, datetime]:
months = []
for path in paths:
match = re.fullmatch(r"yellow_tripdata_(\d{4})-(\d{2})\.parquet", path.name)
if match:
year = int(match.group(1))
month = int(match.group(2))
months.append((year, month))
assert months, "Не удалось определить месяцы из имён yellow_tripdata_*.parquet"
start_year, start_month = min(months)
end_year, end_month = max(months)
expected_start = datetime(start_year, start_month, 1)
if end_month == 12:
expected_end = datetime(end_year + 1, 1, 1)
else:
expected_end = datetime(end_year, end_month + 1, 1)
return expected_start, expected_end
print("Spark и MinIO-клиент готовы")
print(f"RAW URI: {RAW_URI}")
In [ ]:
assert DATA_DIR.exists(), (
"Каталог /opt/data/nyc_taxi не найден. "
"Сначала скачай data bundle по инструкции из START_HERE.md и убедись, что ./data смонтирован в Jupyter как /opt/data."
)
bundle_files = sorted(path for path in DATA_DIR.iterdir() if path.is_file())
yellow_files = sorted(DATA_DIR.glob("yellow_tripdata_*.parquet"))
zone_lookup_file = DATA_DIR / "taxi_zone_lookup.csv"
assert yellow_files, (
"Не найдены файлы yellow_tripdata_*.parquet в /opt/data/nyc_taxi. "
"Скачай каноническое подмножество из START_HERE.md."
)
assert zone_lookup_file.exists(), (
"Не найден taxi_zone_lookup.csv. Он нужен и для этого модуля, и для самостоятельного задания. "
"Скачай его по инструкции из START_HERE.md."
)
bundle_report = pd.DataFrame(
[
{
"file": path.name,
"size_bytes": path.stat().st_size,
"size_human": format_bytes(path.stat().st_size),
}
for path in bundle_files
]
).sort_values("file").reset_index(drop=True)
bundle_report
In [ ]:
upload_candidates = [*yellow_files, zone_lookup_file]
upload_report = []
for path in upload_candidates:
object_key = f"{RAW_PREFIX}{path.name}"
local_size = path.stat().st_size
existed_before = object_exists(BUCKET_NAME, object_key)
if existed_before:
metadata = s3.head_object(Bucket=BUCKET_NAME, Key=object_key)
remote_size = metadata["ContentLength"]
if remote_size != local_size:
raise ValueError(
"Raw object already exists with different size: "
f"{object_key} (local={local_size}, remote={remote_size}). "
"Остановись и проверь data bundle вместо перезаписи raw-слоя."
)
action = "skipped_same_size"
else:
s3.upload_file(str(path), BUCKET_NAME, object_key)
metadata = s3.head_object(Bucket=BUCKET_NAME, Key=object_key)
remote_size = metadata["ContentLength"]
action = "uploaded"
upload_report.append(
{
"file": path.name,
"raw_key": object_key,
"action": action,
"local_size_human": format_bytes(local_size),
"remote_size_human": format_bytes(remote_size),
}
)
pd.DataFrame(upload_report)
In [ ]:
raw_objects = list_objects(BUCKET_NAME, RAW_PREFIX)
pd.DataFrame(
[
{
"key": obj["Key"],
"size_human": format_bytes(obj["Size"]),
}
for obj in raw_objects
]
)
In [ ]:
single_file_uri = f"{RAW_URI}yellow_tripdata_2024-01.parquet"
single_month_df = spark.read.parquet(single_file_uri)
print(f"Читаем один raw-файл: {single_file_uri}")
print(f"Строк в одном файле: {single_month_df.count():,}")
In [ ]:
yellow_df = spark.read.parquet(f"{RAW_URI}yellow_tripdata_*.parquet")
row_count = yellow_df.count()
print(f"Прочитано строк: {row_count:,}")
print(f"Количество parquet-файлов в наборе: {len(yellow_files)}")
In [ ]:
yellow_df.show(10, truncate=False)
In [ ]:
input_files = sorted(yellow_df.inputFiles())
pd.DataFrame({"input_file": input_files})
In [ ]:
yellow_df.printSchema()
In [ ]:
schema_df = pd.DataFrame(yellow_df.dtypes, columns=["column", "spark_type"])
print(f"Количество колонок: {len(yellow_df.columns)}")
schema_df
In [ ]:
null_counts_row = yellow_df.select(
*[F.sum(F.col(column).isNull().cast("int")).alias(column) for column in yellow_df.columns]
).toPandas().iloc[0]
null_report = pd.DataFrame(
{
"column": yellow_df.columns,
"null_count": [int(null_counts_row[column]) for column in yellow_df.columns],
}
)
null_report["null_share"] = (null_report["null_count"] / row_count).round(6)
null_report.sort_values(["null_count", "column"], ascending=[False, True]).reset_index(drop=True)
In [ ]:
numeric_columns = [
column
for column, spark_type in yellow_df.dtypes
if spark_type.startswith(("tinyint", "smallint", "int", "bigint", "float", "double", "decimal"))
]
yellow_df.select(*numeric_columns).describe().toPandas()
In [ ]:
required_columns = {
"fare_amount",
"trip_distance",
"total_amount",
"tpep_pickup_datetime",
"PULocationID",
"DOLocationID",
}
missing_columns = sorted(required_columns.difference(yellow_df.columns))
assert not missing_columns, (
"В raw-схеме не хватает ожидаемых колонок: "
f"{missing_columns}. Возможно, источник поменял формат, и модуль требует обновления."
)
expected_start, expected_end = infer_month_bounds(yellow_files)
anomaly_report = pd.DataFrame(
[
{
"check": "fare_amount < 0",
"count": yellow_df.filter(F.col("fare_amount") < 0).count(),
},
{
"check": "trip_distance = 0 and total_amount != 0",
"count": yellow_df.filter(
(F.col("trip_distance") == 0) & (F.col("total_amount") != 0)
).count(),
},
{
"check": "total_amount > 1000",
"count": yellow_df.filter(F.col("total_amount") > 1000).count(),
},
{
"check": f"pickup outside [{expected_start.date()}, {expected_end.date()})",
"count": yellow_df.filter(
(F.col("tpep_pickup_datetime") < F.lit(expected_start))
| (F.col("tpep_pickup_datetime") >= F.lit(expected_end))
).count(),
},
]
)
anomaly_report
In [ ]:
suspicious_rows = yellow_df.filter(
(F.col("fare_amount") < 0)
| ((F.col("trip_distance") == 0) & (F.col("total_amount") != 0))
| (F.col("total_amount") > 1000)
| (F.col("tpep_pickup_datetime") < F.lit(expected_start))
| (F.col("tpep_pickup_datetime") >= F.lit(expected_end))
)
suspicious_rows.select(
"tpep_pickup_datetime",
"passenger_count",
"trip_distance",
"fare_amount",
"total_amount",
"PULocationID",
"DOLocationID",
).show(20, truncate=False)
In [ ]:
# Ваш код: прочитай taxi_zone_lookup.csv из raw-зоны через Spark
In [ ]:
# Ваш код: выведи схему и первые 10 строк
In [ ]:
# Ваш код: посчитай количество зон по Borough и top-5 pickup locations с join на lookup
In [ ]:
# Ваш код: найди ещё одну аномалию или интересный паттерн в данных
In [ ]:
spark.stop()