Блок 6 - S3, HDFS и Apache Iceberg

20 вопросов: S3 vs HDFS, S3A коммиттеры, архитектура Iceberg, CoW vs MoR, компакция, CDC и сравнение с Delta Lake.

core optimization

Почему операция 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':

  1. Spark читает Manifest List (один маленький файл)
  2. Spark читает только нужные Manifest Files (фильтр по partition)
  3. Из манифестов фильтрует data-файлы по min/max статистике колонок
  4. Читает только отобранные файлы

Результат: 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: предполагает отсутствие конфликтов, обнаруживает их при коммите.

Процесс:

  1. Writer A читает текущий Metadata File (snapshot S1)
  2. Writer B параллельно тоже читает snapshot S1
  3. Writer A коммитит - создаёт новый Metadata File (snapshot S2), атомарно обновляет указатель в Catalog
  4. Writer B пытается закоммитить - обнаруживает, что current snapshot изменился (S1 → S2)
  5. Iceberg проверяет изоляцию: файлы, которые изменил Writer B, не пересекаются с файлами Writer A?
  6. Нет конфликта (разные партиции) → успешный коммит
  7. Конфликт (те же файлы) → 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.