Files
mini-lakehouse-lab/notebooks/04_bronze_with_iceberg.ipynb
T
ddadminandClaude Opus 4.6 68c530c855 refactor(course): перенесена загрузка taxi_zone_lookup из домашки в демо-секцию
- Зачем:
  - демо-ячейки последующих модулей не должны зависеть от домашних заданий предыдущих — модуль 5 использует lookup в демо-коде для JOIN.
- Что:
  - модуль 4: секция 8 из домашки (5 пустых ячеек) превращена в демо с заполненным кодом (чтение CSV, CTAS).
  - модуль 4: новая домашка (секция 9) — осмотр metadata lookup-таблицы, сравнение с основной.
  - модуль 4: секции перенумерованы (9→10, 10→11), завершение упоминает обе bronze-таблицы.
  - модуль 5: assert-сообщение для lookup обновлено (убрана ссылка на «самостоятельное задание»).
  - планы модулей 4 и 5: обновлены дизайн-решения, структура секций и раздел рисков.
- Проверка:
  - прогнать готовые ячейки модуля 4 → lookup-таблица создаётся автоматически → модуль 5 проходит без ошибок.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-07 23:47:05 +03:00

32 KiB
Raw Blame History

Модуль 4. Первая рабочая Iceberg-таблица и слой bronze

В этом модуле мы переходим от raw-файлов в MinIO к первой управляемой таблице Iceberg.

Что делаем на практике:

  • проверяем, что raw-данные из Модуля 3 доступны в MinIO;
  • создаём namespace lakehouse.bronze;
  • загружаем raw Yellow Taxi в Iceberg-таблицу через CTAS;
  • сравниваем чтение по физическому пути и чтение через каталог;
  • смотрим snapshots, history, files и физическую структуру таблицы в MinIO.

После этого ноутбука ты должен понимать, почему bronze-слой это уже не «просто папка с parquet», а управляемая таблица с каталогом и metadata.

0. Перед стартом

Перед выполнением ноутбука:

  • подними стенд по START_HERE.md;
  • пройди notebooks/03_raw_ingest_and_first_read.ipynb;
  • убедись, что raw-файлы уже лежат в s3a://lakehouse/raw/nyc_taxi/.

Связка с привычным PostgreSQL / Greenplum:

Классический DWH (Greenplum / PostgreSQL) Lakehouse
Raw-данные Загружены в stg-таблицу или внешний staging Файлы в object storage (raw/)
Первый управляемый слой Явный DDL + INSERT INTO ... SELECT из stg CREATE TABLE ... USING iceberg AS SELECT из raw (CTAS)
Метаданные pg_catalog внутри СУБД JDBC catalog + metadata files в MinIO
История изменений Нет встроенной истории таблицы snapshots, manifests, metadata JSON

Здесь важно различать два уровня:

  • raw: исходные файлы, которые можно перечитать заново;
  • bronze: первая управляемая таблица, которую удобно читать из Spark и позже из Trino.
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}")

1. Проверка raw перед загрузкой в bronze

Сначала убеждаемся, что raw-данные из Модуля 3 действительно доступны в MinIO. Если их нет, строить bronze просто не из чего.

Ключевое различие:

  • raw читается по физическому пути;
  • bronze будет читаться уже как логическая таблица через каталог.
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

2. Чтение raw-файлов через Spark

Сейчас мы ещё работаем с raw напрямую: spark.read.parquet("s3a://...").

Это не таблица Iceberg, а просто набор parquet-файлов по физическому пути. Нам нужно посмотреть на схему и подготовить явный список колонок для CTAS, чтобы не полагаться на SELECT *.

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"])

Ниже список колонок, который мы явно зафиксируем в CTAS.

Почему это лучше, чем SELECT *:

  • порядок и состав колонок становятся явными;
  • проще заметить неожиданные изменения источника;
  • загрузка в bronze остаётся сознательным решением, а не непрозрачным автоматизмом.
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],
})

3. Создание bronze namespace и Iceberg-таблицы

Сначала создаём namespace lakehouse.bronze, а затем выполняем CTAS (CREATE OR REPLACE TABLE ... USING iceberg AS SELECT ...).

Что такое namespace в Iceberg:

  • это логический контейнер для таблиц;
  • если ты работал с PostgreSQL — это ближайший аналог schema внутри базы;
  • запись lakehouse.bronze.nyc_taxi_yellow читается как catalog.namespace.table.

Почему здесь используем именно bronze, а не module_04:

  • имя сразу отражает роль слоя в пайплайне raw -> bronze -> silver;
  • Модуль 5 будет читать именно lakehouse.bronze.*, и семантическое имя делает маршрут прозрачным;
  • в отличие от демонстрационного module_02, этот слой не временный и будет использоваться в следующих модулях.

Это удобный учебный компромисс:

  • один шаг создаёт таблицу и сразу наполняет её;
  • повторный запуск секции остаётся идемпотентным;
  • минус: CREATE OR REPLACE пересоздаёт таблицу и сбрасывает старую историю snapshot-ов.

Для первого bronze-слоя это допустимо. В Модуле 7 мы отдельно разберём, почему в рабочих сценариях нужно аккуратнее относиться к истории таблицы.

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)

Справочная форма DDL, более похожая на привычный PostgreSQL-стиль, могла бы выглядеть так:

CREATE TABLE lakehouse.bronze.nyc_taxi_yellow (
    VendorID BIGINT,
    tpep_pickup_datetime TIMESTAMP,
    tpep_dropoff_datetime TIMESTAMP,
    passenger_count DOUBLE,
    trip_distance DOUBLE,
    RatecodeID DOUBLE,
    store_and_fwd_flag STRING,
    PULocationID BIGINT,
    DOLocationID BIGINT,
    payment_type BIGINT,
    fare_amount DOUBLE,
    extra DOUBLE,
    mta_tax DOUBLE,
    tip_amount DOUBLE,
    tolls_amount DOUBLE,
    improvement_surcharge DOUBLE,
    total_amount DOUBLE,
    congestion_surcharge DOUBLE,
    Airport_fee DOUBLE
) USING iceberg;

Но в этом модуле мы используем именно CTAS, потому что хотим одной операцией и создать таблицу, и загрузить в неё raw-срез.

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"

4. Читаем уже не raw, а управляемую таблицу

Теперь ключевое отличие:

  • spark.read.parquet("s3a://...") читает набор файлов по физическому пути;
  • spark.table("lakehouse.bronze.nyc_taxi_yellow") читает логическую таблицу через каталог Iceberg.

По результату оба варианта могут вернуть похожие строки. Но второй вариант работает с каталогом и metadata: знает схему, историю изменений и согласованное состояние таблицы.

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

Эту разницу полезно зафиксировать в одной таблице:

Свойство Raw-файлы Iceberg-таблица
Каталогизация Нет Да (PostgreSQL JDBC catalog)
История изменений Нет Snapshots
Контроль схемы Выводится из файлов Зафиксирована в metadata
Чтение из Trino Невозможно без внешней таблицы Да
Атомарность записи Нет Да

5. Что знает Iceberg о таблице: snapshots, history, files

У Iceberg есть встроенные metadata-таблицы. Это один из самых наглядных ответов на вопрос, почему таблица не сводится к «папке с parquet-файлами».

Нас интересуют три представления:

  • .snapshots — какие snapshots есть у таблицы;
  • .history — как менялось текущее состояние таблицы;
  • .files — какие data files входят в текущий snapshot.
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()

Как читать поля в .snapshots:

  • committed_at — когда snapshot стал текущим;
  • snapshot_id — уникальный идентификатор версии таблицы;
  • operation — какая операция создала snapshot, например append или replace;
  • summary — краткое описание результата операции: сколько файлов добавлено, сколько строк записано и т.д.

На практике это можно воспринимать как именованный checkpoint таблицы: каждая запись в истории фиксирует согласованное состояние, к которому движок может обратиться.

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

Iceberg хранит не только путь к данным, но и состояние таблицы как набора metadata-артефактов.

Практический смысл:

  • snapshot фиксирует согласованную версию таблицы;
  • manifests описывают, какие data files входят в snapshot;
  • каталог указывает на актуальный metadata-файл таблицы.

Именно поэтому удалить один parquet-файл «руками» из MinIO опаснее, чем кажется: metadata продолжит ссылаться на него, а таблица станет неконсистентной.

6. Физическая структура таблицы в MinIO

Теперь посмотрим на ту же таблицу с физической стороны. Нас интересуют две группы объектов:

  • data/ — parquet-файлы с данными;
  • metadata/ — metadata JSON, manifest list и manifests.

Это хороший момент, чтобы связать логическую таблицу и физическую структуру в object storage.

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))

Обычно после первого CTAS ты увидишь примерно такую картину:

warehouse/bronze/nyc_taxi_yellow/
    - data
        - *.parquet
    - metadata
        - v1.metadata.json
        - snap-*.avro
        - *.avro

Смысл структуры:

  • data/ хранит содержимое таблицы;
  • metadata/ хранит описание состояния таблицы;
  • JDBC catalog знает, где лежит актуальный metadata-файл.

Это и есть практическая разница между «таблица как логический объект» и «папка с parquet-файлами».

7. Почему raw -> bronze похоже на stg -> ods

В классическом DWH часто есть переход stg -> ods: сначала источник попадает в staging, потом появляется первая рабочая таблица, с которой удобно работать дальше.

В этом курсе аналог выглядит так:

  • raw — неизменяемый архив исходных файлов в object storage;
  • bronze — первая управляемая таблица на Iceberg, которую можно читать через каталог, осматривать через metadata-таблицы и позже запрашивать из Trino;
  • silver — следующий слой, где появятся нормализация, очистка и правила трансформаций.

Поэтому bronze не равно raw:

  • raw нужен как воспроизводимая точка входа;
  • bronze нужен как управляемая таблица с каталогом и snapshot-ами;
  • bronze можно пересоздать из raw, если логика меняется.

8. Вторая bronze-таблица: taxi_zone_lookup

В raw-зоне кроме Yellow Taxi лежит ещё один файл — taxi_zone_lookup.csv. Это справочник зон: каждому LocationID соответствует название района (Borough) и зоны (Zone).

В Модуле 5 мы будем обогащать silver-таблицу через LEFT JOIN с этим справочником. Сейчас загрузим его в bronze по тому же паттерну CTAS, но из CSV вместо Parquet.

In [ ]:
raw_zone_df = (
    spark.read
    .option("header", "true")
    .option("inferSchema", "true")
    .csv(ZONE_LOOKUP_RAW_URI)
)

print(f"Строк в raw lookup: {raw_zone_df.count()}")
raw_zone_df.printSchema()
raw_zone_df.show(5, truncate=False)
In [ ]:
raw_zone_df.createOrReplaceTempView("raw_zone_lookup")

ctas_lookup_sql = f"""
CREATE OR REPLACE TABLE {BRONZE_LOOKUP_TABLE}
USING iceberg
AS
SELECT
    CAST(LocationID AS INT) AS LocationID,
    Borough,
    Zone,
    service_zone
FROM raw_zone_lookup
"""

spark.sql(ctas_lookup_sql)

lookup_df = spark.table(BRONZE_LOOKUP_TABLE)
lookup_count = lookup_df.count()

print(f"Строк в bronze lookup: {lookup_count}")
lookup_df.show(5, truncate=False)

9. Самостоятельное задание

Посмотри на внутреннюю структуру только что созданной taxi_zone_lookup — используй те же инструменты, что в Секциях 5 и 6.

  1. Выведи snapshots и files для lakehouse.bronze.taxi_zone_lookup.
  2. Сравни: сколько data files у lookup-таблицы и сколько у nyc_taxi_yellow? Почему такая разница?
In [ ]:
# Ваш код: snapshots и files для taxi_zone_lookup
In [ ]:
# Ваш код: сравни количество data files в lookup и yellow taxi

10. Checkpoint

Проверь себя:

  1. Чем отличается spark.read.parquet("s3a://...") от spark.table("lakehouse.bronze.nyc_taxi_yellow")?
  2. Что хранит snapshot и зачем он нужен?
  3. Где физически лежат данные таблицы и где лежат её метаданные?
  4. Почему bronze это не то же самое, что raw?
  5. Что произойдёт, если удалить один data file из MinIO, но metadata не поменять?
  6. Можешь ли ты показать таблицу и через Spark, и через MinIO Console?

11. Завершение

Мы намеренно не удаляем bronze-таблицы (nyc_taxi_yellow и taxi_zone_lookup). Они понадобятся в Модуле 5, где мы будем строить silver-слой.

In [ ]:
spark.stop()