Блок 6 - S3, HDFS и Apache Iceberg
20 вопросов: S3 vs HDFS, S3A коммиттеры, архитектура Iceberg, CoW vs MoR, компакция, CDC и сравнение с Delta Lake.
Почему операция rename в HDFS дешёвая, а в S3 - дорогая?¶
В HDFS rename - атомарная операция на уровне NameNode: меняется только запись в дереве метаданных. Данные физически не перемещаются. Сложность: O(1).
В S3 нет rename как операции. Переименование - это: (1) скопировать объект с новым ключом, (2) удалить старый. Для директории с N файлами: N операций copy + N операций delete. Сложность: O(n). Это критично для Spark: стандартный FileOutputCommitter использует rename для атомарного коммита результата.
Решение - S3A Committer (см. следующий вопрос): записывает файлы сразу по финальному пути, избегая rename.
Разница между S3A Committers: Directory Staging, Partitioned Staging и Magic¶
| Committer | Как работает | Когда использовать |
|---|---|---|
| Directory Staging | Файлы пишутся во временную директорию, при коммите - один batch rename | Простые случаи, небольшое число файлов |
| Partitioned Staging | Временная директория per-Task → rename per-partition при коммите | Много партиций, параллельный коммит |
| Magic Committer | Файлы пишутся сразу по финальному пути с magic-метками; при коммите просто удаляются метки | Рекомендуется: нет rename, работает на S3 Native |
# Magic Committer - рекомендуемый вариант для S3
spark.conf.set("spark.hadoop.fs.s3a.committer.name", "magic")
spark.conf.set("spark.hadoop.fs.s3a.committer.magic.enabled", "true")
spark.conf.set("spark.sql.sources.commitProtocolClass",
"org.apache.spark.internal.io.cloud.PathOutputCommitProtocol")
Без специального committer'а Spark использует FileOutputCommitter v1 с двойным rename - катастрофически медленно на S3 при большом числе файлов.
Почему Data Locality теряет смысл при работе Spark с S3?¶
Data Locality - запуск Task на том узле, где физически хранятся данные, чтобы избежать сетевой передачи.
В HDFS каждый блок данных хранится на конкретных DataNode'ах - Spark знает, где данные, и запускает Task рядом.
В S3 данные хранятся в объектном хранилище за пределами кластера: нет понятия "данные на узле X". Все Executor'ы читают S3 по сети с одинаковой задержкой. Spark получает уровень локальности ANY для всех Task'ов.
Практическое следствие:
# Для S3-workload отключить ожидание локальности - нет смысла ждать
spark.conf.set("spark.locality.wait", "0")
Как размер блока HDFS влияет на количество партиций в Spark?¶
Spark по умолчанию создаёт одну партицию на один HDFS-блок. Стандартный размер блока HDFS - 128 МБ.
Файл 10 ГБ / 128 МБ блок = ~80 партиций (Tasks при чтении)
Настройка через spark.sql.files.maxPartitionBytes (default: 128 МБ) - позволяет читать несколько блоков в одну партицию:
# Уменьшить число партиций - читать по 256 МБ в одну партицию
spark.conf.set("spark.sql.files.maxPartitionBytes", "268435456") # 256 MB
# Влияет на начальное число партиций при чтении Parquet/ORC с HDFS
df = spark.read.parquet("hdfs://namenode/data/")
print(df.rdd.getNumPartitions()) # ≈ totalBytes / maxPartitionBytes
На S3 аналогично: maxPartitionBytes контролирует разбивку файлов на Tasks.
Как Iceberg решает проблему File Listing при миллионах файлов?¶
В Hive/стандартных Parquet-таблицах Spark должен сделать LIST запрос к файловой системе, чтобы найти все нужные файлы - при миллионах файлов это занимает минуты и перегружает NameNode/S3.
Iceberg хранит каталог файлов в Manifest Files - Parquet-файлах с информацией о каждом data-файле и его min/max статистиках по колонкам.
Iceberg Catalog → Metadata File → Manifest List → Manifest Files → Data Files
При запросе WHERE date = '2024-01-01':
- Spark читает Manifest List (один маленький файл)
- Spark читает только нужные Manifest Files (фильтр по partition)
- Из манифестов фильтрует data-файлы по min/max статистике колонок
- Читает только отобранные файлы
Результат: listing занимает секунды вместо минут, не зависит от числа файлов.
Опишите структуру метаданных Iceberg (Catalog → Data Files)¶
- Catalog - точка входа, хранит указатель на текущий Metadata File таблицы
- Metadata File - JSON с текущей схемой, partition spec, списком снапшотов
- Manifest List - Avro-файл со списком всех Manifest Files данного снапшота + partition summary
- Manifest File - Avro-файл со списком data files, их размерами и col-level min/max статистиками
- Data Files - реальные Parquet/ORC файлы с данными
Copy-on-Write (CoW) vs Merge-on-Read (MoR): подробное сравнение¶
| CoW | MoR | |
|---|---|---|
| UPDATE / DELETE | Перезаписывает весь data file | Пишет небольшой delete-file |
| Чтение | Быстрое (один проход по файлам) | Медленнее (merge data + delete files) |
| Запись | Медленная (copy файла) | Быстрая |
| Compaction | Не нужна | Требуется регулярно |
| Оптимально для | OLAP: много чтений, редкие обновления | CDC: частые точечные обновления |
# Выбор режима при создании таблицы Iceberg
spark.sql("""
CREATE TABLE catalog.events (id LONG, ts TIMESTAMP, data STRING)
USING iceberg
TBLPROPERTIES (
'write.delete.mode' = 'merge-on-read', -- или 'copy-on-write'
'write.update.mode' = 'merge-on-read',
'write.merge.mode' = 'merge-on-read'
)
""")
MoR в Iceberg: при DELETE создаётся Equality Delete File (содержит ключи удалённых строк) или Positional Delete File (содержит файл + позиция строки). При чтении Spark выполняет merge "на лету".
Как Iceberg обеспечивает Snapshot Isolation?¶
Каждая запись в Iceberg создаёт новый снапшот - новый Manifest List указывает на новый набор файлов. Читатели работают со старым снапшотом до момента явного обновления.
# Читатель видит снапшот на момент открытия scan
df = spark.read.table("catalog.events") # snapshot_id зафиксирован в плане
# Параллельный writer создаёт новый снапшот
spark.sql("INSERT INTO catalog.events VALUES (1, now(), 'new')")
# Читатель df всё ещё видит старые данные - изоляция!
# Явный Time Travel на конкретный снапшот
df_old = spark.read.option("snapshot-id", "1234567890").table("catalog.events")
При конкурентной записи Iceberg использует Optimistic Concurrency Control: операция читает текущее состояние, выполняет изменения, и при коммите проверяет, не изменились ли файлы, от которых она зависит. Если конфликт - откат и повтор.
Что такое Hidden Partitioning в Iceberg?¶
В Hive/стандартном Spark пользователь должен явно писать WHERE year = 2024 AND month = 1 чтобы получить partition pruning - иначе Spark читает все партиции. Партиционирование "протекает" в SQL.
В Iceberg партиционирование скрыто: пользователь пишет фильтр по реальной колонке (WHERE ts > '2024-01-01'), а Iceberg автоматически транслирует его в partition pruning.
-- Таблица партиционирована по часам из timestamp-колонки
CREATE TABLE events (ts TIMESTAMP, user_id LONG, event STRING)
USING iceberg
PARTITIONED BY (hours(ts)); -- Скрытое партиционирование
-- Запрос - пользователь фильтрует по ts, не думая о партициях
SELECT * FROM events WHERE ts BETWEEN '2024-01-01' AND '2024-01-02';
-- Iceberg автоматически пропускает нерелевантные часовые партиции
Partition Evolution: можно изменить логику партиционирования (например, с days(ts) на months(ts)) без перезаписи старых данных - Iceberg применяет новую логику только к новым файлам.
Какие изменения схемы Iceberg поддерживает без перезаписи данных?¶
Iceberg поддерживает schema evolution - изменение схемы не требует перезаписи существующих файлов:
| Операция | Поддержка | Что происходит |
|---|---|---|
| Добавить колонку | ✓ | Старые файлы возвращают null для новой колонки |
| Удалить колонку | ✓ | Старые файлы читаются, колонка скрывается |
| Переименовать колонку | ✓ | Маппинг по column ID, не по имени |
| Изменить тип | Частично | int→long, float→double, date→datetime |
| Сменить порядок колонок | ✓ | Маппинг по column ID |
-- Добавить колонку (не перезаписывает данные)
ALTER TABLE catalog.events ADD COLUMN session_id STRING;
-- Переименовать (старые файлы читаются правильно через column ID)
ALTER TABLE catalog.events RENAME COLUMN user_id TO customer_id;
Iceberg использует column IDs (числовые идентификаторы) вместо имён - поэтому переименование безопасно. Hive-таблицы хранят данные по позиции колонки - переименование может поломать старые файлы.
Зачем и когда запускать compaction (rewriteDataFiles) в Iceberg?¶
После многих мелких записей (streaming sink, частые MERGE) таблица накапливает много мелких файлов - каждый требует отдельного S3 запроса, планирование замедляется.
Compaction (rewriteDataFiles) перечитывает файлы и объединяет в более крупные оптимальные файлы:
# Компакция - объединить мелкие файлы в оптимальные 128–512 MB
spark.sql("""
CALL catalog.system.rewrite_data_files(
table => 'events',
strategy => 'sort',
sort_order => 'zorder(user_id, event_date)',
options => map(
'target-file-size-bytes', '134217728', -- 128 MB
'min-file-size-bytes', '33554432' -- перезаписать файлы < 32 MB
)
)
""")
Когда запускать:
- После streaming pipeline (ежечасно/ежедневно)
- После mass-DELETE/UPDATE (файлы с delete-markers в MoR)
- При деградации query performance (много мелких файлов в explain плане)
Как Iceberg разрешает конкурентные записи (Optimistic Concurrency Control)?¶
Iceberg использует Optimistic Concurrency: предполагает отсутствие конфликтов, обнаруживает их при коммите.
Процесс:
- Writer A читает текущий Metadata File (snapshot S1)
- Writer B параллельно тоже читает snapshot S1
- Writer A коммитит - создаёт новый Metadata File (snapshot S2), атомарно обновляет указатель в Catalog
- Writer B пытается закоммитить - обнаруживает, что current snapshot изменился (S1 → S2)
- Iceberg проверяет изоляцию: файлы, которые изменил Writer B, не пересекаются с файлами Writer A?
- Нет конфликта (разные партиции) → успешный коммит
- Конфликт (те же файлы) → retry с новым планом
# Retry логика встроена в Iceberg
# spark.sql("MERGE INTO ...") - автоматически повторяет при конфликте
Для серьёзной конкурентности рекомендуется RESTCatalog (Polaris, Unity) - поддерживает distributed locking через RDBMS.
Как использовать Iceberg как Sink для Spark Structured Streaming?¶
Iceberg поддерживает потоковую запись из Structured Streaming начиная с Iceberg 0.12:
stream_df = spark.readStream \
.format("kafka") \
.option("subscribe", "events") \
.load() \
.selectExpr("CAST(value AS STRING) AS json")
parsed = stream_df.select(from_json(col("json"), schema).alias("d")).select("d.*")
# Запись в Iceberg с checkpoint
query = parsed.writeStream \
.format("iceberg") \
.outputMode("append") \
.option("checkpointLocation", "s3://bucket/checkpoints/stream_events/") \
.trigger(processingTime="1 minute") \
.toTable("catalog.events") # Spark 3.x: прямая запись в каталог
query.awaitTermination()
Особенности:
- Каждый micro-batch = один Iceberg снапшот (транзакционная запись)
- Exactly-once: checkpoint + Iceberg transaction log = нет дублей при рестарте
- Compaction нужно запускать отдельно - streaming создаёт много мелких файлов
Как реализовать инкрементальное чтение из таблицы Iceberg?¶
Iceberg поддерживает чтение только новых снапшотов - данных, добавленных с момента последней обработки:
# Incremental Read - только данные между двумя снапшотами
df = spark.read \
.format("iceberg") \
.option("start-snapshot-id", "1234567890") \ # включительно или нет
.option("end-snapshot-id", "9876543210") \
.load("catalog.events")
# Или через Spark SQL
df = spark.sql("""
SELECT * FROM catalog.events
VERSION BETWEEN 1234567890 AND 9876543210
""")
# Для Structured Streaming - читать только новые снапшоты
stream = spark.readStream \
.format("iceberg") \
.option("stream-from-timestamp", "2024-01-01T00:00:00") \
.load("catalog.events")
Практический паттерн: хранить last_processed_snapshot_id в отдельной таблице метаданных, передавать его как start-snapshot-id при следующем запуске.
Iceberg vs Delta Lake vs Apache Hudi: основные различия¶
| Аспект | Apache Iceberg | Delta Lake | Apache Hudi |
|---|---|---|---|
| Метаданные | Avro manifest files | JSON transaction log | Timeline (Avro) |
| Engine-agnostic | ✓ (Spark, Flink, Trino, Presto) | Частично (основной: Spark) | Частично |
| CDC / MoR | ✓ (Equality + Positional delete files) | ✓ (Deletion vectors) | ✓ (основная фича) |
| Hidden Partitioning | ✓ | ✗ | ✗ |
| Partition Evolution | ✓ без перезаписи | ✗ | ✗ |
| Schema Evolution | Расширенная (column IDs) | Базовая | Базовая |
| Governance | Apache Software Foundation | Linux Foundation (Delta.io) | Apache Software Foundation |
| Лицензия | Apache 2.0 | Apache 2.0 | Apache 2.0 |
Когда выбирать Iceberg: open ecosystem (несколько движков), сложная partition evolution, enterprise multi-tenant.
Когда выбирать Delta Lake: Databricks-окружение, простота начала работы, тесная интеграция с Databricks.
Когда выбирать Hudi: высокая частота обновлений (CDC), уже используется AWS EMR.
Как настроить Spark для борьбы с S3 throttling (503 Slow Down)?¶
S3 ограничивает пропускную способность по prefix'у: ~3500 PUT/s и ~5500 GET/s на один prefix. При параллельной записи Spark легко превышает эти лимиты.
# Retry параметры S3A
spark.conf.set("spark.hadoop.fs.s3a.retry.limit", "20")
spark.conf.set("spark.hadoop.fs.s3a.retry.interval", "500ms")
spark.conf.set("spark.hadoop.fs.s3a.retry.throttle.interval", "1s")
# Увеличить connection pool
spark.conf.set("spark.hadoop.fs.s3a.connection.maximum", "200")
spark.conf.set("spark.hadoop.fs.s3a.threads.max", "64")
# Multipart upload для больших файлов (снижает число запросов)
spark.conf.set("spark.hadoop.fs.s3a.multipart.size", "128M")
spark.conf.set("spark.hadoop.fs.s3a.fast.upload", "true")
Архитектурное решение: prefix sharding - использовать разные prefix'ы (папки) для разных партиций, чтобы распределить нагрузку по нескольким S3 "shard'ам". Iceberg делает это автоматически через hash-based file naming.
Что такое Bloom Filters в Iceberg и как их использовать?¶
Bloom Filter - компактная вероятностная структура данных, позволяющая быстро проверить, есть ли конкретное значение в файле. False positive возможен, false negative - нет.
Iceberg хранит Bloom Filter в метаданных каждого Parquet-файла. При запросе WHERE user_id = 'abc' Spark проверяет Bloom Filter файла перед чтением - если "точно нет" → пропускает файл.
-- Включить Bloom Filter для конкретной колонки при создании таблицы
CREATE TABLE catalog.events (id LONG, user_id STRING, ...)
USING iceberg
TBLPROPERTIES (
'write.parquet.bloom-filter-enabled.column.user_id' = 'true',
'write.parquet.bloom-filter-max-bytes' = '1048576' -- 1 MB per file
);
Эффективен для:
- High-cardinality колонок (user_id, order_id, session_id)
- Точечных поисков (
WHERE id = 'xxx') - Lookup по нескольким значениям (
WHERE id IN (...))
Не помогает: range-фильтры (WHERE amount > 100) - для них min/max статистики Iceberg эффективнее.
Как S3 решил проблему Eventual Consistency?¶
До декабря 2020 года S3 давал Eventual Consistency для некоторых операций:
- PUT объекта → немедленно виден при GET по тому же ключу
- но LIST после PUT мог не показывать новый объект несколько секунд
- DELETE + PUT с тем же ключом → GET мог вернуть старый объект
Это было критично для Spark: Task записывает файл, другой процесс делает LIST и не видит его → неполный результат.
Решения до 2020: AWS S3Guard (DynamoDB как consistent metadata store), Hadoop S3A fs.s3a.consistent=true, специальные committer'ы.
С декабря 2020: AWS объявила S3 Strong Consistency - все операции (PUT, DELETE, LIST, GET) теперь строго согласованы. Spark на AWS S3 больше не требует S3Guard или специальных workaround'ов для consistency.
# Сейчас не нужны workaround'ы для S3 consistency:
# spark.conf.set("fs.s3a.consistent", "true") # устарело
# spark.conf.set("fs.s3a.metadatastore.impl", "org.apache.hadoop.fs.s3a.s3guard.DynamoDBMetadataStore") # устарело
Для S3-совместимых хранилищ (MinIO, Ceph) - гарантии зависят от реализации: MinIO предоставляет strong consistency, Ceph - в зависимости от конфигурации.
Разница между типами каталогов Iceberg: HadoopCatalog, HiveCatalog, RESTCatalog¶
| Catalog | Хранение метаданных | Требования | Когда использовать |
|---|---|---|---|
| HadoopCatalog | Файловая система (HDFS/S3) - путь к metadata файлам | Только доступ к FS | Dev/test, простые single-team случаи |
| HiveCatalog | Hive Metastore (Thrift API) | HMS running | Существующий Hive/Spark кластер |
| RESTCatalog | HTTP REST API (Polaris, Unity, Nessie) | REST-сервер | Production multi-engine, multi-tenant |
# HadoopCatalog - простейший, без сервера
spark.conf.set("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog")
spark.conf.set("spark.sql.catalog.local.type", "hadoop")
spark.conf.set("spark.sql.catalog.local.warehouse", "s3a://bucket/warehouse/")
# HiveCatalog - интеграция с Hive Metastore
spark.conf.set("spark.sql.catalog.hive_prod", "org.apache.iceberg.spark.SparkCatalog")
spark.conf.set("spark.sql.catalog.hive_prod.type", "hive")
spark.conf.set("spark.sql.catalog.hive_prod.uri", "thrift://metastore:9083")
# RESTCatalog - production (Polaris, Apache Gravitino, Nessie)
spark.conf.set("spark.sql.catalog.rest_prod", "org.apache.iceberg.spark.SparkCatalog")
spark.conf.set("spark.sql.catalog.rest_prod.catalog-impl",
"org.apache.iceberg.rest.RESTCatalog")
spark.conf.set("spark.sql.catalog.rest_prod.uri", "https://catalog.example.com/")
RESTCatalog - рекомендуется для production: поддерживает multi-engine (Spark + Flink + Trino + Hive одновременно), distributed locking для concurrent writes, RBAC.
Как управлять историей снапшотов: expire snapshots и delete orphan files¶
После многих операций таблица накапливает старые снапшоты и orphan-файлы (файлы, на которые никто не ссылается). Без очистки они занимают место на S3/HDFS.
Expire Snapshots - удаляет старые снапшоты (и соответствующие metadata/manifest файлы):
-- SQL: удалить снапшоты старше 7 дней, оставить минимум 5
CALL catalog.system.expire_snapshots(
table => 'events',
older_than => TIMESTAMP '2024-01-01 00:00:00',
retain_last => 5
);
# Python/PySpark
from pyiceberg.catalog import load_catalog
catalog = load_catalog("rest", uri="https://catalog.example.com/")
table = catalog.load_table("db.events")
table.expire_snapshots().expire_older_than(
datetime(2024, 1, 1, tzinfo=timezone.utc)
).retain_last(5).commit()
Delete Orphan Files - находит и удаляет файлы, которые не упоминаются ни в одном снапшоте (могут появиться после сбоя задачи записи):
CALL catalog.system.remove_orphan_files(
table => 'events',
older_than => TIMESTAMP '2024-01-01 00:00:00'
);
Рекомендуемое расписание: expire_snapshots - ежедневно, remove_orphan_files - еженедельно. Не удалять снапшоты, нужные для Time Travel - использовать retain_last.