Каталог и иерархия метаданных
Apache Iceberg — это открытая спецификация табличного формата для аналитических данных на object storage (S3, GCS, ADLS, HDFS). В отличие от Delta Lake, где единственным источником истины является _delta_log/ внутри директории таблицы, Iceberg использует catalog-first архитектуру: внешний каталог хранит указатель на текущий metadata file, а вся остальная информация образует дерево иммутабельных файлов.
В Модуле 11 мы разобрали линейный transaction log Delta Lake. Iceberg идёт другим путём — трёхуровневая иерархия метаданных обеспечивает O(1) доступ к текущему состоянию таблицы без replay журнала.
Этот курс использует engine-agnostic подход. Примеры кода — на Python с библиотекой pyiceberg (текущий стабильный релиз 0.11.1, март 2026), без привязки к Apache Spark. Iceberg — это спецификация, а не движок. Текущая версия спецификации — V2, с V3 для deletion vectors.
Catalog-first vs Path-based
Delta Lake — path-based: достаточно знать путь к директории таблицы (s3://bucket/warehouse/my_table/), чтобы прочитать _delta_log/ и восстановить состояние. Каталог (Hive Metastore, Unity Catalog) — опционален.
Iceberg — catalog-first: чтобы найти таблицу, нужно обратиться к каталогу. Каталог возвращает путь к текущему metadata file. Без каталога таблицу не открыть.
Delta Lake (path-based)
Delta Lake: клиент напрямую обращается к файловой системе по известному пути. Каталог опционален — можно работать без него.Клиент → Путь → _delta_log/
Клиент читает _delta_log/ директорию, replay-ит JSON-коммиты, находит checkpoint — и получает текущее состояние таблицы.Iceberg (catalog-first)
Iceberg: клиент обязательно обращается к каталогу (REST, Hive, Glue, JDBC). Каталог хранит единственный указатель — путь к текущему metadata file.Клиент → Каталог → Metadata file
Каталог возвращает location текущего metadata file. Клиент читает этот JSON-файл — и сразу получает полное состояние таблицы. Никакого replay.Почему catalog-first?
| Преимущество | Механизм |
|---|---|
| Atomic table updates | Каталог атомарно переключает указатель с одного metadata file на другой (compare-and-swap) |
| Кросс-движковая совместимость | Spark, Flink, Trino, DuckDB, Dremio — все обращаются к одному каталогу и видят одинаковое состояние |
| Безопасный rename | Переименование таблицы — изменение маппинга в каталоге, без копирования данных |
| Multi-table transactions | Каталог может координировать атомарные изменения нескольких таблиц (в рамках REST API v2) |
Iceberg можно открыть без каталога через StaticTable (напрямую по пути к metadata file), но это read-only режим без гарантий консистентности. В продакшене всегда используйте каталог.
Трёхуровневая иерархия метаданных
Каталог хранит единственное значение — путь к текущему metadata file. Дальше разворачивается дерево:
Каталог (catalog)
Каталог (REST, Hive, Glue, JDBC) — единственная мутабельная точка во всей системе. Хранит маппинг: имя таблицы → путь к текущему metadata file. Атомарно переключается при каждом коммите.Уровень 1: Metadata file (JSON)
Metadata file — JSON-файл с полным описанием таблицы: все схемы за всю историю, все partition specs, все snapshots, текущий snapshot ID, свойства таблицы, sort orders. Один файл = полное состояние.Уровень 2: Manifest list (Avro)
Manifest list — Avro-файл, содержащий список manifest files для конкретного snapshot. Каждая запись: путь к manifest file, partition spec ID, число added/existing/deleted файлов, range partition значений. Позволяет пропускать irrelevant manifests без чтения.Уровень 3: Manifest file (Avro)
Manifest file — Avro-файл, содержащий список data files с их метаданными: путь, формат, record count, file size bytes, column-level min/max/null_count/distinct_count/NaN_count. Позволяет partition pruning и data skipping на уровне отдельных файлов.Ключевое свойство: всё, кроме каталога, иммутабельно. Metadata files, manifest lists, manifest files, data files — создаются один раз и никогда не модифицируются. Это обеспечивает snapshot isolation, time travel и безопасное конкурентное чтение.
Уровень 1: Metadata file (JSON)
Metadata file — это JSON-файл (обычно с расширением .metadata.json), содержащий полное описание таблицы. Каждый коммит в таблицу создаёт новый metadata file.
s3://warehouse/db/orders/metadata/
├── v1.metadata.json ← CREATE TABLE
├── v2.metadata.json ← первый INSERT
├── v3.metadata.json ← ALTER TABLE ADD COLUMN
├── v4.metadata.json ← UPDATE
└── v5.metadata.json ← текущий (каталог указывает сюда)
Структура metadata file
{
"format-version": 2,
"table-uuid": "af23c9d7-fff1-4a5a-a2c8-55c59bd782aa",
"location": "s3://warehouse/db/orders",
"last-sequence-number": 5,
"last-updated-ms": 1700000500000,
"last-column-id": 7,
"current-schema-id": 1,
"schemas": [ ... ],
"default-spec-id": 0,
"partition-specs": [ ... ],
"last-partition-id": 1000,
"default-sort-order-id": 0,
"sort-orders": [ ... ],
"properties": {
"write.parquet.compression-codec": "zstd",
"commit.retry.num-retries": "4"
},
"current-snapshot-id": 3497810934857103984,
"snapshots": [ ... ],
"snapshot-log": [ ... ],
"metadata-log": [ ... ],
"refs": { ... }
}
Подробный разбор полей:
schemas — история схем
{
"schemas": [
{
"schema-id": 0,
"type": "struct",
"fields": [
{ "id": 1, "name": "order_id", "type": "long", "required": true },
{ "id": 2, "name": "customer_id", "type": "long", "required": true },
{ "id": 3, "name": "amount", "type": "decimal(10,2)", "required": true },
{ "id": 4, "name": "order_date", "type": "timestamptz", "required": true }
]
},
{
"schema-id": 1,
"type": "struct",
"fields": [
{ "id": 1, "name": "order_id", "type": "long", "required": true },
{ "id": 2, "name": "customer_id", "type": "long", "required": true },
{ "id": 3, "name": "amount", "type": "decimal(10,2)", "required": true },
{ "id": 4, "name": "order_date", "type": "timestamptz", "required": true },
{ "id": 5, "name": "status", "type": "string", "required": false },
{ "id": 6, "name": "region", "type": "string", "required": false }
]
}
],
"current-schema-id": 1
}
Ключевые свойства:
- Уникальные
field-id: каждое поле получает уникальный ID при создании. Полеstatus(id=5) всегда остаётся id=5, даже если другие поля удаляются. ID никогда не переиспользуются. - Полная история: все предыдущие схемы сохраняются. Старые snapshot-ы ссылаются на старые schema-id.
- Привязка по
field-id, не по имени: Parquet-файлы, записанные со старой схемой, корректно читаются с новой — сопоставление по field-id, а не по имени или позиции колонки.
Именно привязка по field-id делает schema evolution в Iceberg надёжнее, чем в Delta Lake (привязка по имени) и в Parquet (привязка по позиции). Переименование колонки не ломает чтение старых данных.
snapshots — массив снимков
{
"snapshots": [
{
"snapshot-id": 3497810934857103984,
"parent-snapshot-id": 2894710293847102938,
"sequence-number": 5,
"timestamp-ms": 1700000500000,
"manifest-list": "s3://warehouse/db/orders/metadata/snap-3497810934857103984-0-abc123.avro",
"summary": {
"operation": "append",
"added-data-files": "3",
"added-records": "150000",
"total-data-files": "42",
"total-records": "5200000"
},
"schema-id": 1
}
]
}
Поле manifest-list — ссылка на следующий уровень иерархии.
refs — ветки и теги
{
"refs": {
"main": {
"snapshot-id": 3497810934857103984,
"type": "branch",
"max-ref-age-ms": null
},
"audit-2025-01": {
"snapshot-id": 2894710293847102938,
"type": "tag",
"max-ref-age-ms": 7776000000
}
}
}
- Branch — мутабельная ссылка, перемещается при каждом коммите (аналог
HEADв git). По умолчанию —main. - Tag — иммутабельная ссылка, фиксирует конкретный snapshot (аналог git tag). Используется для аудита, compliance, фиксации версий данных.
Уровень 2: Manifest list (Avro)
Manifest list — Avro-файл, содержащий список manifest files для конкретного snapshot. Каждая запись в manifest list описывает один manifest file:
Manifest pruning на уровне manifest list
Поле partitions — это механизм manifest pruning: если запрос фильтрует по order_date >= '2025-01-01', а manifest list показывает, что manifest #3 содержит данные только за декабрь 2024, этот manifest пропускается без скачивания. Для таблицы с тысячами manifest files это экономит секунды на планировании.
from pyiceberg.catalog import load_catalog
catalog = load_catalog("default", type="sql",
uri="sqlite:////tmp/warehouse/catalog.db",
warehouse="file:///tmp/warehouse")
table = catalog.load_table("db.orders")
# Метаданные текущего snapshot
snapshot = table.current_snapshot()
print(f"Snapshot ID: {snapshot.snapshot_id}")
print(f"Sequence number: {snapshot.sequence_number}")
print(f"Manifest list: {snapshot.manifest_list}")
# Список manifest files
manifests = snapshot.manifests(table.io)
for m in manifests:
print(f" {m.manifest_path}")
print(f" added_files: {m.added_files_count}")
print(f" existing_files: {m.existing_files_count}")
print(f" deleted_files: {m.deleted_files_count}")
print(f" partition_spec_id: {m.partition_spec_id}")
Уровень 3: Manifest file (Avro)
Manifest file — Avro-файл, содержащий записи о data files. Каждая запись — это data_file struct:
Data skipping с column-level statistics
Поля lower_bounds и upper_bounds — основа data skipping. Для каждого файла Iceberg хранит min/max значения по каждой колонке. При выполнении запроса с фильтром движок сравнивает предикат с границами:
-- Запрос: найти заказы дороже $1000 за январь 2025
SELECT * FROM orders
WHERE amount > 1000
AND order_date >= '2025-01-01'
AND order_date < '2025-02-01'
Движок проверяет каждый data file в manifest:
upper_bound(amount) <= 1000→ skip (в файле нет записей с amount > 1000)upper_bound(order_date) < '2025-01-01'→ skip (файл не содержит январских данных)lower_bound(order_date) >= '2025-02-01'→ skip (файл содержит только февральские+ данные)
Только файлы, прошедшие все проверки, отправляются на чтение. Для широких таблиц с десятками колонок data skipping может пропустить 95%+ файлов.
В отличие от Delta Lake, где min/max хранятся в JSON-строке внутри add action (и парсятся двухуровнево), Iceberg хранит границы в типизированных Avro-полях с binary-сериализацией. Это быстрее при десктопах тысяч файлов.
Чтение через pyiceberg
from pyiceberg.catalog import load_catalog
catalog = load_catalog("default", type="sql",
uri="sqlite:////tmp/warehouse/catalog.db",
warehouse="file:///tmp/warehouse")
table = catalog.load_table("db.orders")
# Полная метаинформация
print(f"Table UUID: {table.metadata.table_uuid}")
print(f"Format version: {table.metadata.format_version}")
print(f"Location: {table.metadata.location}")
print(f"Current schema: {table.schema()}")
print(f"Partition spec: {table.spec()}")
# Все data files через plan_files()
scan = table.scan()
for task in scan.plan_files():
df = task.file
print(f" {df.file_path}")
print(f" format: {df.file_format}")
print(f" records: {df.record_count}")
print(f" size: {df.file_size_in_bytes / 1024 / 1024:.1f} MB")
# Чтение данных в Arrow
arrow_table = table.scan(
row_filter="amount > 1000",
selected_fields=("order_id", "customer_id", "amount")
).to_arrow()
print(f"Filtered rows: {len(arrow_table)}")
Типы каталогов
Iceberg спецификация определяет интерфейс каталога, но не реализацию. Существует несколько реализаций:
REST Catalog — стандарт де-факто
REST Catalog — стандартизированный HTTP API, определённый в Iceberg REST OpenAPI spec. Ключевые эндпоинты:
| Endpoint | Метод | Описание |
|---|---|---|
/v1/namespaces | GET | Список namespace-ов |
/v1/namespaces/{ns}/tables | GET | Список таблиц в namespace |
/v1/namespaces/{ns}/tables/{table} | GET | Метаданные таблицы (metadata location) |
/v1/namespaces/{ns}/tables/{table} | POST | Создание таблицы |
/v1/namespaces/{ns}/tables/{table} | POST | Commit (atomic update) |
/v1/config | GET | Конфигурация каталога |
/v1/oauth/tokens | POST | OAuth2 аутентификация |
Преимущество REST каталога — credential vending: каталог может выдавать временные credentials для доступа к storage (STS tokens для S3, signed URLs для GCS). Клиенту не нужны долгоживущие credentials к storage — достаточно авторизоваться в каталоге.
Настройка pyiceberg с разными каталогами
from pyiceberg.catalog import load_catalog
# REST Catalog (рекомендовано для продакшена)
rest_catalog = load_catalog("production", **{
"type": "rest",
"uri": "https://polaris.example.com/api/catalog",
"credential": "client_id:client_secret",
"warehouse": "my_warehouse",
"scope": "PRINCIPAL_ROLE:ALL",
})
# AWS Glue
glue_catalog = load_catalog("aws", **{
"type": "glue",
"s3.region": "us-east-1",
})
# JDBC (PostgreSQL)
jdbc_catalog = load_catalog("jdbc", **{
"type": "sql",
"uri": "postgresql://user:pass@host:5432/iceberg_catalog",
"warehouse": "s3://my-bucket/warehouse",
})
# SQL Catalog (локальная разработка)
local_catalog = load_catalog("local", **{
"type": "sql",
"uri": "sqlite:////tmp/warehouse/catalog.db",
"warehouse": "file:///tmp/warehouse",
})
# Hive Metastore
hive_catalog = load_catalog("hive", **{
"type": "hive",
"uri": "thrift://hive-metastore:9083",
"warehouse": "s3://my-bucket/warehouse",
})
Apache Polaris — reference-реализация REST Catalog
Apache Polaris (incubating, версия 1.3.0) — открытый REST-каталог для Iceberg, инициированный Snowflake и переданный в Apache Incubator:
Spark
Spark — один из движков, подключающихся к Polaris через стандартный Iceberg REST Catalog API. Использует credential vending для доступа к S3/GCS/ADLS.Trino
Trino — аналитический движок, подключающийся к Polaris. Поддерживает federated query через Iceberg REST API.Flink
Flink — streaming-движок. Polaris поддерживает dynamic sink и credential vending для streaming ingestion.pyiceberg
pyiceberg — Python-клиент. Подключается к Polaris напрямую для ad-hoc анализа, schema management, table maintenance.Apache Polaris (REST Catalog Server)
Polaris Server — Quarkus-based сервер, реализующий Iceberg REST OpenAPI. Управляет каталогами, namespace-ами, таблицами, правами доступа. RBAC: principal → principal role → catalog role → privileges.Ключевые возможности Polaris 1.3.0:
- RBAC: principal → principal role → catalog role → privileges (CATALOG_MANAGE_CONTENT, TABLE_READ_DATA, TABLE_WRITE_DATA и др.)
- Credential vending: временные STS tokens для S3, SigV4 аутентификация
- Generic Tables GA: каталогизация таблиц других форматов (Hudi, Delta Lake)
- Multi-engine: Spark, Flink, Trino, StarRocks, Dremio, DuckDB, pyiceberg
Сравнение с Delta Lake
| Аспект | Delta Lake | Apache Iceberg |
|---|---|---|
| Источник истины | _delta_log/ (файлы в директории таблицы) | Каталог → metadata file → manifest tree |
| Каталог | Опционален (Unity Catalog) | Обязателен (REST, Hive, Glue, JDBC) |
| Доступ к состоянию | O(n) replay JSON + checkpoint | O(1) — один metadata file |
| Метаданные data files | JSON-строка stats в add action | Типизированные Avro-поля в manifest file |
| Schema evolution привязка | По имени колонки | По field-id (column ID mapping встроен) |
| Partition evolution | (partition columns фиксированы) | Есть (разные partition specs для разных manifest) |
| Multi-engine | delta-rs, delta-spark | Широкая экосистема (Spark, Flink, Trino, и др.) |
Практика: инспекция метаданных с pyiceberg
from pyiceberg.catalog import load_catalog
import json
# Создаём каталог и таблицу
catalog = load_catalog("demo", type="sql",
uri="sqlite:////tmp/iceberg_demo/catalog.db",
warehouse="file:///tmp/iceberg_demo")
catalog.create_namespace_if_not_exists("analytics")
import pyarrow as pa
schema = pa.schema([
("order_id", pa.int64()),
("customer_id", pa.int64()),
("amount", pa.float64()),
("order_date", pa.timestamp("us", tz="UTC")),
("region", pa.string()),
])
table = catalog.create_table_if_not_exists(
"analytics.orders", schema=schema
)
# Записываем данные
import pyarrow as pa
from datetime import datetime, timezone
data = pa.table({
"order_id": [1, 2, 3, 4, 5],
"customer_id": [100, 101, 100, 102, 101],
"amount": [250.0, 1500.0, 75.0, 3200.0, 890.0],
"order_date": [
datetime(2025, 1, 15, tzinfo=timezone.utc),
datetime(2025, 1, 16, tzinfo=timezone.utc),
datetime(2025, 2, 1, tzinfo=timezone.utc),
datetime(2025, 2, 10, tzinfo=timezone.utc),
datetime(2025, 3, 5, tzinfo=timezone.utc),
],
"region": ["EU", "US", "EU", "APAC", "US"],
})
table.append(data)
# Инспекция иерархии метаданных
meta = table.metadata
print("=== Metadata file ===")
print(f" format-version: {meta.format_version}")
print(f" table-uuid: {meta.table_uuid}")
print(f" location: {meta.location}")
print(f" current-snapshot-id: {meta.current_snapshot_id}")
print(f" schemas: {len(meta.schemas)} versions")
print(f" partition-specs: {len(meta.partition_specs)}")
print("\n=== Current snapshot ===")
snap = table.current_snapshot()
print(f" snapshot-id: {snap.snapshot_id}")
print(f" sequence-number: {snap.sequence_number}")
print(f" summary: {dict(snap.summary)}")
print("\n=== Manifest list → Manifest files ===")
manifests = snap.manifests(table.io)
for i, m in enumerate(manifests):
print(f" Manifest {i}: {m.manifest_path}")
print(f" added_files: {m.added_files_count}")
print(f" content: {'data' if m.content == 0 else 'deletes'}")
Для локальной разработки и тестирования используйте type="sql" с SQLite. Для продакшена — REST каталог (Polaris, Unity Catalog, Gravitino) или AWS Glue. JDBC (PostgreSQL) — хороший self-hosted вариант.
Итоги
- Iceberg = catalog-first: каталог — единственная мутабельная точка, всё остальное иммутабельно
- Три уровня метаданных: metadata file (JSON) → manifest list (Avro) → manifest file (Avro) → data files
- O(1) доступ к состоянию: каталог указывает на текущий metadata file, нет replay журнала
- Column-level statistics: min/max/null_count в manifest file обеспечивают data skipping без чтения данных
- Manifest pruning: partition summary в manifest list позволяет пропускать целые manifest files
- REST Catalog — стандарт: Apache Polaris (1.3.0-incubating) — reference-реализация с RBAC и credential vending
- Schema evolution по field-id: безопаснее привязки по имени (Delta Lake) или по позиции (Parquet)