- Зачем: - нужен практический модуль 4, который переводит студента от raw-файлов к первой управляемой Iceberg-таблице. - Что: - добавлен ноутбук notebooks/04_bronze_with_iceberg.ipynb с CTAS в bronze, разбором metadata tables и структурой хранения в MinIO. - добавлены пояснения про namespace, сравнение raw vs Iceberg и разбор snapshot-полей по итогам ревью. - обновлен plans/README.md строкой для модуля 4. - Проверка: - синтаксис всех code-ячеек ноутбука проверен локальной компиляцией. - тесты пройдены пользователем.
31 KiB
31 KiB
In [ ]:
from pathlib import PurePosixPath
import boto3
import pandas as pd
from pyspark.sql import SparkSession, functions as F
BUCKET_NAME = "lakehouse"
RAW_PREFIX = "raw/nyc_taxi/"
RAW_URI = f"s3a://{BUCKET_NAME}/{RAW_PREFIX}"
RAW_YELLOW_URI = f"{RAW_URI}yellow_tripdata_*.parquet"
ZONE_LOOKUP_RAW_URI = f"{RAW_URI}taxi_zone_lookup.csv"
CATALOG_NAME = "lakehouse"
BRONZE_NAMESPACE = f"{CATALOG_NAME}.bronze"
BRONZE_TABLE = f"{BRONZE_NAMESPACE}.nyc_taxi_yellow"
BRONZE_LOOKUP_TABLE = f"{BRONZE_NAMESPACE}.taxi_zone_lookup"
spark = SparkSession.builder .appName("module-04-bronze-with-iceberg") .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 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 get_table_location(table_name: str) -> str:
details = spark.sql(f"DESCRIBE TABLE EXTENDED {table_name}").toPandas()
location_row = details.loc[details["col_name"] == "Location", "data_type"]
assert not location_row.empty, f"Не удалось найти Location для {table_name}"
return location_row.iloc[0]
print("Spark и MinIO-клиент готовы")
print(f"RAW URI: {RAW_URI}")
print(f"Bronze table: {BRONZE_TABLE}")
In [ ]:
raw_objects = list_objects(BUCKET_NAME, RAW_PREFIX)
assert raw_objects, (
"Raw-зона пуста. Сначала пройди Модуль 3 и загрузи data bundle в MinIO."
)
raw_report = pd.DataFrame(
[
{
"key": obj["Key"],
"size_bytes": obj["Size"],
"size_human": format_bytes(obj["Size"]),
}
for obj in raw_objects
]
).sort_values("key").reset_index(drop=True)
yellow_object_count = int(raw_report["key"].str.contains(r"yellow_tripdata_.*\.parquet", regex=True).sum())
assert yellow_object_count > 0, (
"В raw-зоне не найдены yellow_tripdata_*.parquet. Сначала пройди Модуль 3."
)
assert (raw_report["key"] == f"{RAW_PREFIX}taxi_zone_lookup.csv").any(), (
"В raw-зоне не найден taxi_zone_lookup.csv. Он нужен для задания этого модуля и join в Модуле 5."
)
print(f"Всего объектов в raw-зоне: {len(raw_report)}")
print(f"Parquet-файлов Yellow Taxi: {yellow_object_count}")
raw_report
In [ ]:
raw_yellow_df = spark.read.parquet(RAW_YELLOW_URI)
raw_row_count = raw_yellow_df.count()
print(f"Прочитано raw-строк: {raw_row_count:,}")
print(f"Количество колонок: {len(raw_yellow_df.columns)}")
raw_yellow_df.show(10, truncate=False)
In [ ]:
raw_yellow_df.printSchema()
pd.DataFrame(raw_yellow_df.dtypes, columns=["column", "spark_type"])
In [ ]:
expected_columns = [
"VendorID",
"tpep_pickup_datetime",
"tpep_dropoff_datetime",
"passenger_count",
"trip_distance",
"RatecodeID",
"store_and_fwd_flag",
"PULocationID",
"DOLocationID",
"payment_type",
"fare_amount",
"extra",
"mta_tax",
"tip_amount",
"tolls_amount",
"improvement_surcharge",
"total_amount",
"congestion_surcharge",
"Airport_fee",
]
missing_columns = sorted(set(expected_columns).difference(raw_yellow_df.columns))
assert not missing_columns, (
"В raw-схеме не хватает ожидаемых колонок: "
f"{missing_columns}. Проверь data bundle и инструкцию из START_HERE.md."
)
pd.DataFrame({
"column": expected_columns,
"present_in_raw": [column in raw_yellow_df.columns for column in expected_columns],
})
In [ ]:
spark.sql(f"CREATE NAMESPACE IF NOT EXISTS {BRONZE_NAMESPACE}")
raw_yellow_df.createOrReplaceTempView("raw_yellow")
ctas_sql = f"""
CREATE OR REPLACE TABLE {BRONZE_TABLE}
USING iceberg
AS
SELECT
VendorID,
tpep_pickup_datetime,
tpep_dropoff_datetime,
passenger_count,
trip_distance,
RatecodeID,
store_and_fwd_flag,
PULocationID,
DOLocationID,
payment_type,
fare_amount,
extra,
mta_tax,
tip_amount,
tolls_amount,
improvement_surcharge,
total_amount,
congestion_surcharge,
Airport_fee
FROM raw_yellow
"""
print("Temp view ready: raw_yellow")
print(ctas_sql)
In [ ]:
# Для полного набора из нескольких месяцев эта ячейка может выполняться 1-3 минуты.
spark.sql(ctas_sql)
bronze_df = spark.table(BRONZE_TABLE)
bronze_row_count = bronze_df.count()
print(f"Строк в raw: {raw_row_count:,}")
print(f"Строк в bronze: {bronze_row_count:,}")
assert bronze_row_count == raw_row_count, "Количество строк в bronze не совпало с raw"
In [ ]:
bronze_df.show(10, truncate=False)
In [ ]:
comparison_df = pd.DataFrame(
[
{
"mode": "raw path",
"reader": f"spark.read.parquet('{RAW_YELLOW_URI}')",
"row_count": raw_row_count,
"object": "набор parquet-файлов",
},
{
"mode": "catalog table",
"reader": f"spark.table('{BRONZE_TABLE}')",
"row_count": bronze_row_count,
"object": "Iceberg-таблица",
},
]
)
comparison_df
In [ ]:
spark.sql(
f"SELECT committed_at, snapshot_id, parent_id, operation, manifest_list, summary "
f"FROM {BRONZE_TABLE}.snapshots ORDER BY committed_at DESC"
).toPandas()
In [ ]:
spark.sql(
f"SELECT made_current_at, snapshot_id, parent_id, is_current_ancestor "
f"FROM {BRONZE_TABLE}.history ORDER BY made_current_at DESC"
).toPandas()
In [ ]:
files_df = spark.sql(
f"""
SELECT
content,
file_path,
file_format,
record_count,
file_size_in_bytes
FROM {BRONZE_TABLE}.files
ORDER BY file_path
"""
)
files_pd = files_df.toPandas()
files_pd["file_size_human"] = files_pd["file_size_in_bytes"].map(format_bytes)
files_pd
In [ ]:
file_summary = pd.DataFrame(
[
{
"metric": "data_files_count",
"value": len(files_pd),
},
{
"metric": "total_records_in_files",
"value": int(files_pd["record_count"].sum()) if not files_pd.empty else 0,
},
{
"metric": "total_file_size_bytes",
"value": int(files_pd["file_size_in_bytes"].sum()) if not files_pd.empty else 0,
},
{
"metric": "total_file_size_human",
"value": format_bytes(int(files_pd["file_size_in_bytes"].sum())) if not files_pd.empty else "0 B",
},
]
)
file_summary
In [ ]:
bronze_location = get_table_location(BRONZE_TABLE)
bronze_prefix = bronze_location.replace(f"s3a://{BUCKET_NAME}/", "").rstrip("/") + "/"
bronze_objects = list_objects(BUCKET_NAME, bronze_prefix)
assert bronze_objects, f"По пути {bronze_location} не найдены объекты таблицы"
bronze_storage_report = pd.DataFrame(
[
{
"key": obj["Key"],
"kind": "DATA" if "/data/" in obj["Key"] else "META" if "/metadata/" in obj["Key"] else "OTHER",
"size_bytes": obj["Size"],
"size_human": format_bytes(obj["Size"]),
}
for obj in bronze_objects
]
).sort_values(["kind", "key"]).reset_index(drop=True)
print(f"Table location: {bronze_location}")
bronze_storage_report
In [ ]:
storage_summary = bronze_storage_report.groupby("kind", as_index=False)["size_bytes"].sum()
storage_summary["size_human"] = storage_summary["size_bytes"].map(format_bytes)
storage_summary
In [ ]:
root = PurePosixPath(bronze_prefix.rstrip("/"))
rendered_lines = [str(root)]
for key in bronze_storage_report["key"]:
rel = PurePosixPath(key).relative_to(root)
parts = rel.parts
for depth, part in enumerate(parts):
prefix = " " * depth + ("- " if depth else "")
rendered = prefix + part
if rendered not in rendered_lines:
rendered_lines.append(rendered)
print("Структура каталога таблицы:\n")
print("\n".join(rendered_lines))
In [ ]:
# Ваш код: прочитай taxi_zone_lookup.csv из raw-зоны через Spark
In [ ]:
# Ваш код: создай таблицу lakehouse.bronze.taxi_zone_lookup
In [ ]:
# Ваш код: загрузи данные в bronze lookup-таблицу
In [ ]:
# Ваш код: проверь результат через spark.table(...)
In [ ]:
# Ваш код: посмотри snapshots для lakehouse.bronze.taxi_zone_lookup
In [ ]:
spark.stop()