Rename Problem: почему S3 не POSIX и как это ломает Spark

Rename Problem - главный кошмар инженеров, переходящих с HDFS на облачные хранилища. Разбираем физику S3, механику FileOutputCommitter, S3A Committers и как Lakehouse-форматы полностью решили проблему.

storage

Почему это важно знать каждому Spark-инженеру

Когда команда переезжает с on-premise Hadoop кластера на облачный Data Lake поверх S3 или S3-совместимого хранилища (MinIO, Ceph, Yandex Object Storage), первое время всё выглядит привычно: те же пути вида /bucket/path/to/data/, те же партиции, та же команда spark.write.parquet(...). Но очень быстро начинаются странности:

  • Задачи записи, которые на HDFS занимали 10 минут, на S3 длятся 2–3 часа.
  • Драйвер Spark зависает на стадии «commit» и не отвечает на сигналы.
  • После падения задачи в output-директории оказываются частичные данные - как будто запись «прошла наполовину».
  • Перезапуск задачи создаёт дубликаты строк в «перезаписанных» партициях.

Всё это - проявления одной фундаментальной проблемы: Rename Problem. Spark исторически спроектирован под POSIX-семантику файловых систем, а S3 - это не файловая система. Непонимание этого различия приводит к потере данных, деградации производительности в 10–100 раз и огромным счетам за облачный compute.

Этот урок - детальный разбор физики проблемы, её последствий и современных решений.

Объектное хранилище против файловой системы: фундаментальные различия

Иллюзия иерархии в S3

Когда вы видите путь s3://my-bucket/data/year=2024/month=06/part-00000.parquet, у вас возникает ощущение привычной директорной структуры: папка data, внутри year=2024, внутри month=06, файл part-00000.parquet. Это иллюзия.

В S3 не существует директорий как таких. Внутри бакета хранятся объекты с уникальными ключами. Ключ - это просто строка. Символ / - это обычный символ, такой же как любой другой. Вся «иерархия», которую вы видите в UI AWS Console или в выводе aws s3 ls, - это визуальная проекция, которую строит клиент, группируя ключи по общим префиксам.

Физическая реальность S3:
┌─────────────────────────────────────────────────────────────┐
│  Бакет: my-bucket                                           │
│  Тип:   Flat key-value store                                │
│                                                             │
│  Ключ → Значение (blob bytes):                              │
│  "data/year=2024/month=06/part-00000.parquet" → [bytes...]  │
│  "data/year=2024/month=06/part-00001.parquet" → [bytes...]  │
│  "data/year=2024/month=07/part-00000.parquet" → [bytes...]  │
│  "README.md"                                  → [bytes...]  │
└─────────────────────────────────────────────────────────────┘

Что видит клиент (иллюзия директорий):
data/
  year=2024/
    month=06/
      part-00000.parquet
      part-00001.parquet
    month=07/
      part-00000.parquet
README.md

Эта разница кажется тривиальной, но именно она порождает всю цепочку проблем.

POSIX-семантика: что ожидает Spark

POSIX (Portable Operating System Interface) - это стандарт, описывающий поведение операционных систем, включая файловые системы. Spark, как и большинство Java/JVM-фреймворков, работает через абстракцию FileSystem из Hadoop, которая исторически строилась с расчётом на POSIX-совместимые системы.

Ключевые POSIX-гарантии, на которые опирается Spark:

1. Атомарное переименование (rename). В POSIX rename(oldpath, newpath) - атомарная операция. Она либо полностью выполняется, либо нет. Нельзя наблюдать промежуточное состояние. Если несколько процессов одновременно делают rename в один и тот же путь, ровно один из них «победит». На HDFS rename реализован через изменение метаданных в NameNode - это действительно атомарная O(1) операция независимо от размера данных.

2. Консистентное листинг директорий. После создания файла он немедленно виден при ls той же директории. Не «в течение нескольких секунд», не «вероятно» - немедленно и всегда.

3. Иерархические директории. Создание файла /a/b/c/file.txt подразумевает существование директорий /a, /a/b, /a/b/c. Директории - это реальные объекты с метаданными, временем изменения и правами доступа.

4. Блокировка файлов (flock). Процесс может заблокировать файл для эксклюзивного доступа. Другие процессы будут ждать.

Анатомия операции rename в S3

Когда Hadoop FileSystem реализует rename для S3, она вынуждена делать это в несколько шагов, каждый из которых - отдельный HTTP API-вызов к S3:

  1. Листинг источника: ListObjectsV2 - получить список всех объектов под исходным префиксом. Для директории из 100 000 файлов это занимает ~10–50 секунд (S3 отдаёт максимум 1000 объектов за запрос).

  2. Копирование каждого объекта: CopyObject для каждого файла по отдельности. Каждый вызов - HTTP-запрос с латентностью ~20–100 мс. Данные при этом перемещаются внутри S3 (server-side copy), но при большом размере объектов может потребоваться Multipart Copy.

  3. Удаление старых объектов: DeleteObjects (batch до 1000 ключей за раз) - удалить все исходные объекты.

Rename директории с 10 000 файлами по 100 MB:

ListObjectsV2:   10 000 / 1000 = 10 запросов × 50ms = 500ms
CopyObject:      10 000 × 100ms = 16 минут (при параллельности 1)
                 10 000 × 100ms / 100 потоков = ~10 секунд
                 + пропускная способность: 1 TB данных через S3
DeleteObjects:   10 000 / 1000 = 10 запросов × 20ms = 200ms

Итого: от 10 секунд до 16+ минут в зависимости от параллельности
Сравните с rename в HDFS: ~50ms (изменение метаданных NameNode)

Критически важно: эта операция не атомарна. Если процесс упадёт в момент, когда скопировано 7000 из 10 000 файлов, в хранилище будут существовать одновременно старые и новые пути. Нет никакого механизма «rollback».

Под капотом Spark: как работает FileOutputCommitter

Зачем вообще нужен коммитер

Представьте, что Spark пишет датафрейм из 100 партиций параллельно на 50 экзекуторах. Каждый экзекутор пишет свою часть данных в выходную директорию. Проблемы:

  • Экзекутор может упасть в середине записи - файл будет частичным.
  • Spark может перезапустить таск на другом экзекуторе - тот же файл будет записан дважды (speculative execution).
  • Если джоба целиком упадёт - в output-директории окажутся мусорные файлы.

FileOutputCommitter - это механизм, гарантирующий, что читатели видят либо полный результат всей джобы, либо не видят ничего. Это реализация принципа «всё или ничего» для distributed batch writes.

Algorithm v1: двойной rename

Алгоритм v1 (исторически первый) работает в два этапа переименования:

На HDFS: каждый rename - это ~50 мс вызов NameNode. 1000 тасков = 1000 × 50 мс = 50 секунд. Приемлемо.

На S3: каждый rename - это Copy + Delete для каждого файла в директории таска. При 2 файлах на таск и 1000 тасков - 2000 CopyObject запросов, потом 2000 DeleteObjects. При больших файлах и медленной сети - часы работы.

Хуже того, v1 делает двойной rename: сначала attempt → task, потом task → output. На S3 это O(N²) с точки зрения overhead: данные копируются дважды через сеть S3 API.

Algorithm v2: быстрее, но опаснее

v2 упрощает схему: убирает промежуточный _temporary/task_XXX/ слой. Экзекуторы пишут напрямую в _temporary/attempt_XXX/, а job commit переименовывает оттуда прямо в output:

v1: attempt_dir → task_dir → output/  (2 rename)
v2: attempt_dir → output/             (1 rename)

На S3 v2 быстрее, но создаёт критическую проблему атомарности. Во время job commit (финального переименования) идёт серия CopyObject запросов. Если процесс падает в середине - часть файлов уже в output/, часть ещё в _temporary/. Читатели видят частичный результат. Data corruption.

Настройка алгоритма

# Выбор алгоритма в spark-submit или SparkConf:
spark = SparkSession.builder \
    .config("spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version", "2") \
    .getOrCreate()

# ❌ v1 на S3 = катастрофически медленно (двойное копирование)
# ❌ v2 на S3 = быстрее, но data corruption при падении

# Ни один из вариантов не является правильным для production S3!

Реальные последствия: что именно ломается

Сценарий 1: зависание драйвера при job commit

Spark-джоба записывает 500 партиций по 1 GB каждая. Все 500 экзекуторов успешно завершили запись в _temporary/. Начинается job commit на драйвере.

С алгоритмом v1 драйвер должен переименовать 500 директорий. Каждое переименование для 1 GB файла - это CopyObject (1 GB данных), потом DeleteObject. S3 имеет максимум 5 GB/s пропускной способности на аккаунт. 500 GB = ~100 секунд только на IO. Но при последовательном выполнении (как делает стандартный коммитер) - это 500 × ~10 секунд = 83 минуты. Задача, которая на HDFS коммитилась за 30 секунд, на S3 коммитится почти 1,5 часа.

Симптомы в Spark UI:
- Все 500 тасков: SUCCEEDED (100%)
- Stage status: "Commit in progress"
- Driver logs: "Committing output of task attempt_0..."
- Прошло 2 часа, статус не меняется
- Heartbeat timeout → драйвер помечается мёртвым
- Джоба помечается как FAILED

Сценарий 2: data pollution при overwrite

Spark пишет обновлённую партицию year=2024/month=06/ с алгоритмом v2. Задача:

  1. Удалить старую партицию.
  2. Записать новые файлы.

FileOutputCommitter v2 при overwrite делает:

  1. Пишет новые данные в _temporary/.
  2. Начинает job commit: переименовывает файлы из _temporary/ в year=2024/month=06/.
  3. Падает на середине (OutOfMemoryError на драйвере, network timeout, etc).

Результат: в year=2024/month=06/ содержится смесь новых и старых файлов. Следующая джоба читает это как корректные данные. Тихая data corruption.

Сценарий 3: дубликаты при speculative execution

Spark запускает speculative copy таска на втором экзекуторе (оригинальный таск медленнее среднего). Оба экзекутора пишут один и тот же партиционный файл. Первый завершает запись, коммитит. Второй тоже завершает и тоже пытается скоммитить.

На HDFS: второй rename провалится, потому что rename атомарен - нельзя переименовать поверх существующего. На S3: оба CopyObject успевают завершиться, оба файла записаны. В output - дубликаты строк.

S3A Committers: решение от Hadoop AWS

Чтобы устранить проблему rename, Hadoop проект разработал специализированные коммитеры для S3, которые полностью исключают операцию rename из протокола записи. Они входят в модуль hadoop-aws начиная с версии 3.1.

Magic Committer: мультипарт-загрузка как транзакция

Magic Committer использует механизм S3 Multipart Upload (MPU) как замену атомарному rename.

Как работает S3 Multipart Upload: S3 позволяет загружать большой объект по частям:

  1. CreateMultipartUpload - создать «транзакцию» загрузки, получить UploadId.
  2. UploadPart × N - загрузить части данных. Части сами по себе невидимы для других читателей.
  3. CompleteMultipartUpload - атомарно «применить» все части, создав объект. Только после этого объект становится видимым.
  4. AbortMultipartUpload - отменить загрузку, удалить загруженные части.

Magic Committer использует это как транзакционный механизм:

Ключевое преимущество: Job Commit делает только CompleteMultipartUpload запросы - они мгновенны (это metadata-операция в S3). Данные уже загружены в S3, просто «невидимы». Нет CopyObject, нет переноса байт, нет зависания.

Почему «Magic»: коммитер использует специальные «магические» пути в S3, где незавершённые multipart uploads хранятся невидимыми до commit. Это выглядит магически с точки зрения файловой системы.

# Конфигурация Magic Committer:
spark = SparkSession.builder \
    .config("spark.hadoop.fs.s3a.committer.name", "magic") \
    .config("spark.hadoop.fs.s3a.bucket.all.committer.magic.enabled", "true") \
    .config("spark.sql.sources.commitProtocolClass",
            "org.apache.spark.internal.io.cloud.PathOutputCommitProtocol") \
    .config("spark.sql.parquet.output.committer.class",
            "org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter") \
    .getOrCreate()

Ограничение Magic Committer: S3 накапливает незавершённые multipart uploads, которые тарифицируются (хотя и дёшево). Нужно настроить lifecycle rule для автоматического удаления abandoned MPU:

{
  "Rules": [{
    "ID": "abort-incomplete-multipart",
    "Status": "Enabled",
    "Filter": {"Prefix": ""},
    "AbortIncompleteMultipartUpload": {"DaysAfterInitiation": 7}
  }]
}

Staging Committer: локальный диск как буфер

Staging Committer работает по-другому: экзекуторы пишут данные не сразу в S3, а в локальную файловую систему (локальный диск узла). Финальный commit выгружает данные с локальных дисков напрямую в финальные пути S3 через PUT Object.

Преимущества Staging Committer:

  • Task commits используют локальную POSIX FS - атомарны, мгновенны.
  • Данные в S3 пишутся сразу в финальный путь через PUT Object (не CopyObject).
  • Нет промежуточных _temporary/ директорий в S3.

Недостатки Staging Committer:

  • Требует достаточно места на локальных дисках экзекуторов (временное хранилище = размер batch).
  • Job commit всё ещё занимает время - это PUT в S3, хотя и без overhead копирования.
# Конфигурация Directory Staging Committer:
spark = SparkSession.builder \
    .config("spark.hadoop.fs.s3a.committer.name", "directory") \
    .config("spark.hadoop.fs.s3a.committer.staging.tmp.path", "/tmp/spark-staging") \
    .config("spark.sql.sources.commitProtocolClass",
            "org.apache.spark.internal.io.cloud.PathOutputCommitProtocol") \
    .getOrCreate()

# Partitioned Staging Committer (для партиционированных таблиц):
spark = SparkSession.builder \
    .config("spark.hadoop.fs.s3a.committer.name", "partitioned") \
    .config("spark.hadoop.fs.s3a.committer.staging.conflict-mode", "replace") \
    .getOrCreate()
# conflict-mode: fail | append | replace

Сравнение коммитеров

Параметр FileOutputCommitter v1 FileOutputCommitter v2 Magic Directory Staging
Скорость коммита (S3) ❌ O(N²) данных ❌ O(N) данных ✅ O(1) метаданные ✅ O(N) upload
Атомарность ✅ (через двойной rename) ❌ Нет ✅ MPU ✅ (локальный rename + PUT)
Корректность при retry ❌ Data pollution
Локальное место Нет Нет Нет ⚠️ Нужно
Поддержка S3

Lakehouse-форматы: радикальное решение проблемы

S3A Committers решают проблему rename для «сырых» файловых записей. Но настоящее решение пришло с появлением Lakehouse-форматов: Delta Lake, Apache Iceberg и Apache Hudi. Они полностью переосмысливают протокол записи, делая rename вообще ненужным.

Принцип: метаданные как транзакция

Ключевая идея Lakehouse-форматов: данные и метаданные разделены. Файлы данных пишутся в произвольные пути (часто UUID-based), и только после успешной записи всех файлов атомарно обновляются метаданные таблицы - указатель на новый набор файлов.

Delta Lake: transaction log как источник истины

В Delta Lake таблица - это директория с двумя элементами:

  • _delta_log/ - директория с JSON и Parquet файлами транзакционного лога.
  • Файлы данных - Parquet файлы с UUID-именами.

Каждый commit записывает файл вида _delta_log/0000000000000000042.json. Этот JSON содержит атомарную операцию: добавить такие-то файлы, удалить такие-то файлы, изменить такие-то метаданные.

// _delta_log/0000000000000000042.json
{
  "commitInfo": {"timestamp": 1717660800000, "operation": "WRITE"},
  "add": {
    "path": "year=2024/month=06/part-abc123.parquet",
    "size": 104857600,
    "stats": {"numRecords": 1000000}
  },
  "add": {
    "path": "year=2024/month=06/part-def456.parquet",
    "size": 98765432,
    "stats": {"numRecords": 950000}
  }
}

Атомарность: запись одного JSON-файла в S3 через PUT Object - атомарная операция. Либо файл создан полностью, либо нет. Если Spark падает до записи лог-файла - данные (Parquet) уже лежат в S3, но они «невидимы» для читателей таблицы. Следующий vacuum удалит их как orphan files.

# Запись в Delta Lake:
df.write \
    .format("delta") \
    .mode("append") \                       # или "overwrite"
    .partitionBy("year", "month") \
    .save("s3a://my-bucket/delta/orders")

# Никаких _temporary/ директорий, никаких rename!
# Файлы идут сразу в финальный путь с UUID-именами
# Коммит = запись одного JSON в _delta_log/

Apache Iceberg: snapshot-based атомарность

Iceberg использует более сложную, но более надёжную систему метаданных:

s3://bucket/table/
├── metadata/
│   ├── v1.metadata.json      # Snapshot 1
│   ├── v2.metadata.json      # Snapshot 2
│   ├── snap-123.avro         # Manifest list
│   └── manifest-456.avro     # Manifest file
└── data/
    ├── part-abc.parquet
    └── part-def.parquet

Атомарный коммит в Iceberg - это atomic swap указателя на current metadata:

  1. Записать новые data files в data/.
  2. Записать manifest file с перечнем новых файлов.
  3. Записать manifest list с указателем на новый manifest.
  4. Записать новый vN.metadata.json с новым snapshot.
  5. Атомарно обновить указатель metadata/version-hint.text на новую версию.

Шаг 5 - это атомарный PUT Object в S3. Ни Rename, ни Copy, ни Delete.

# Запись в Iceberg:
df.write \
    .format("iceberg") \
    .mode("append") \
    .save("catalog.schema.orders")

# Или через SQL:
spark.sql("""
    INSERT INTO catalog.schema.orders
    SELECT * FROM staging_orders
""")

Как Lakehouse-форматы решают все проблемы rename

Eventual Consistency: историческая проблема и её решение

До декабря 2020 года S3 имел eventual consistency: после записи файла, его могло не быть видно в ListObjects несколько секунд. После удаления - файл мог ещё секунду отдаваться при GET. Это порождало сложные race conditions в Spark.

Что происходило:

  1. Executor записывал файл → получал 200 OK от S3.
  2. Другой экзекутор делал ListObjects → файл ещё не видел.
  3. Spark driver считал, что файл не записан → ошибка или пропуск данных.

AWS решила проблему в декабре 2020 года: S3 теперь гарантирует strong read-after-write consistency для всех операций (GET, HEAD, LIST) после любой успешной записи. Это устранило целый класс проблем для AWS S3.

Но осторожно: On-premise S3-совместимые хранилища (MinIO, Ceph RGW, Yandex Object Storage) могут не иметь таких гарантий или иметь их с оговорками. При переходе на self-hosted object storage - проверяйте consistency semantics вашего продукта.

# Исторический workaround (больше не нужен для AWS S3):
# s3guard - DynamoDB как consistent metadata store
spark.conf.set("fs.s3a.metadatastore.impl",
               "org.apache.hadoop.fs.s3a.s3guard.DynamoDBMetadataStore")
# После декабря 2020: AWS официально убрала s3guard как deprecated

S3 Throttling: ещё одна production-проблема

Когда Spark работает с высокой параллельностью (сотни воркеров одновременно), возникает другая проблема: S3 throttling. S3 ограничивает количество запросов по одному префиксу:

  • 3500 PUT/COPY/POST/DELETE запросов в секунду на префикс (partition key в S3).
  • 5500 GET/HEAD запросов в секунду на префикс.

При 200 параллельных экзекуторах, каждый из которых делает по 50 PUT в секунду - это 10 000 PUT/с, что вдвое превышает лимит. S3 начинает возвращать HTTP 503 Slow Down.

Что такое S3 partition key: S3 автоматически партиционирует данные по первым символам ключа (prefix). Если все файлы пишутся в s3://bucket/data/year=2024/month=06/, они попадают в одну «partition» S3, которая обслуживается ограниченным набором серверов.

Решение - salting prefix:

# ❌ Проблема: все файлы в одном префиксе
df.write.parquet("s3a://bucket/events/year=2024/month=06/")
# Все PUT идут на один S3 partition key → throttling

# ✅ Решение 1: рандомизировать структуру префиксов
# Добавить хеш-бакет как первый уровень
df.withColumn("hash_bucket",
    F.pmod(F.hash(F.col("user_id")), F.lit(16)).cast("string")
).write \
    .partitionBy("hash_bucket", "year", "month") \
    .parquet("s3a://bucket/events/")
# Теперь: events/hash_bucket=0/year=2024/month=06/
#          events/hash_bucket=1/year=2024/month=06/
# 16 разных префиксов → нагрузка распределена

# ✅ Решение 2: Iceberg/Delta автоматически используют UUID-файлы
# UUID: a1b2c3d4-e5f6-...
# Разные UUID начинаются с разных символов → автоматически разные S3 partitions

Чеклист: настройка Spark для работы с S3

Полный набор конфигураций для production Spark + S3:

from pyspark.sql import SparkSession

def create_spark_for_s3(
    app_name: str,
    s3_endpoint: str = None,       # None для AWS, URL для MinIO
    committer: str = "magic"       # magic | directory | partitioned
) -> SparkSession:
    """
    Создать SparkSession с оптимальной конфигурацией для S3.
    """
    builder = SparkSession.builder.appName(app_name)

    # --- Базовая аутентификация и endpoint ---
    if s3_endpoint:
        # MinIO или другой S3-совместимый сервис
        builder = builder \
            .config("spark.hadoop.fs.s3a.endpoint", s3_endpoint) \
            .config("spark.hadoop.fs.s3a.path.style.access", "true")

    # --- S3A Committer (вместо FileOutputCommitter) ---
    builder = builder \
        .config("spark.hadoop.fs.s3a.committer.name", committer) \
        .config("spark.sql.sources.commitProtocolClass",
                "org.apache.spark.internal.io.cloud.PathOutputCommitProtocol") \
        .config("spark.sql.parquet.output.committer.class",
                "org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter")

    if committer == "magic":
        builder = builder \
            .config("spark.hadoop.fs.s3a.bucket.all.committer.magic.enabled", "true")

    if committer in ("directory", "partitioned"):
        builder = builder \
            .config("spark.hadoop.fs.s3a.committer.staging.tmp.path",
                    "/tmp/spark-staging") \
            .config("spark.hadoop.fs.s3a.committer.staging.conflict-mode",
                    "replace")  # fail | append | replace

    # --- Производительность I/O ---
    builder = builder \
        .config("spark.hadoop.fs.s3a.connection.maximum", "200") \
        .config("spark.hadoop.fs.s3a.threads.max", "64") \
        .config("spark.hadoop.fs.s3a.multipart.size", "134217728") \       # 128 MB
        .config("spark.hadoop.fs.s3a.fast.upload", "true") \
        .config("spark.hadoop.fs.s3a.fast.upload.buffer", "array") \
        .config("spark.hadoop.fs.s3a.block.size", "134217728") \           # 128 MB
        .config("spark.hadoop.fs.s3a.readahead.range", "2097152")          # 2 MB

    # --- Retry и throttling ---
    builder = builder \
        .config("spark.hadoop.fs.s3a.retry.limit", "7") \
        .config("spark.hadoop.fs.s3a.retry.interval", "500ms") \
        .config("spark.hadoop.fs.s3a.attempts.maximum", "20") \
        .config("spark.hadoop.fs.s3a.connection.timeout", "200000") \
        .config("spark.hadoop.fs.s3a.socket.recv.buffer", "65536")

    # --- Отключить Legacy FileOutputCommitter ---
    builder = builder \
        .config("spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version", "2") \
        .config("spark.hadoop.mapreduce.fileoutputcommitter.cleanup-failures.ignored",
                "true")

    return builder.getOrCreate()


# Для Lakehouse (Delta/Iceberg) - ещё проще, коммитер не нужен:
def create_spark_for_iceberg_s3(app_name: str) -> SparkSession:
    return SparkSession.builder \
        .appName(app_name) \
        .config("spark.sql.extensions",
                "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
        .config("spark.sql.catalog.spark_catalog",
                "org.apache.iceberg.spark.SparkSessionCatalog") \
        .config("spark.sql.catalog.spark_catalog.type", "hive") \
        .config("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog") \
        .config("spark.sql.catalog.local.type", "hadoop") \
        .config("spark.sql.catalog.local.warehouse", "s3a://my-bucket/warehouse") \
        .getOrCreate()

Anti-patterns при работе Spark + S3

Антипаттерн 1: df.coalesce(1).write.csv()

# ❌ Классическая ошибка новичков:
df.coalesce(1) \
  .write \
  .mode("overwrite") \
  .csv("s3a://bucket/output/result.csv")

# Что происходит:
# 1. Весь датафрейм собирается на один экзекутор (OOM risk!)
# 2. Записывается один файл в _temporary/
# 3. Rename: CopyObject (весь файл!) → boto3 таймаут при файлах > 5GB
# 4. Пользователь видит "result.csv/" директорию, а не файл

Антипаттерн 2: прямое использование os.rename() или shutil для S3 путей

# ❌ Использование Python OS API на S3 путях:
import os
os.rename("s3a://bucket/temp/output", "s3a://bucket/final/output")
# → TypeError: не работает с s3:// протоколом
# Или хуже: использование boto3 напрямую без понимания, что это copy+delete

import boto3
s3 = boto3.client("s3")
s3.copy_object(  # "Переименование" папки - нужно делать для каждого объекта!
    CopySource={"Bucket": "my-bucket", "Key": "temp/output/file.parquet"},
    Bucket="my-bucket",
    Key="final/output/file.parquet"
)
s3.delete_object(Bucket="my-bucket", Key="temp/output/file.parquet")
# ↑ Нужно делать для каждого файла в папке - легко ошибиться

Антипаттерн 3: SaveMode.Overwrite без Lakehouse-формата

# ❌ Опасный overwrite без транзакционного лога:
df.write \
    .mode("overwrite") \     # Сначала удаляет все файлы, потом пишет новые
    .parquet("s3a://bucket/data/orders/")

# Проблема: между удалением старых файлов и появлением новых
# читатели видят пустую таблицу!
# При падении пайплайна → потеря всех данных

# ✅ Безопасный overwrite с Delta Lake:
df.write \
    .format("delta") \
    .mode("overwrite") \     # Транзакционный overwrite: старое видно до коммита
    .save("s3a://bucket/delta/orders/")
# Атомарно: читатели видят старые данные до завершения коммита,
# после коммита - новые данные, никакого окна "пустой таблицы"

Антипаттерн 4: множество мелких файлов

# ❌ Запись тысяч мелких файлов:
streaming_df.writeStream \
    .trigger(processingTime="1 second") \  # каждую секунду новые файлы
    .format("parquet") \
    .option("path", "s3a://bucket/events/") \
    .start()

# Результат через сутки:
# 86400 файлов × средний размер 1MB = 86GB данных в 86400 файлах
# Проблема: каждый файл = отдельный HTTP GET запрос при чтении
# Spark читает 86400 файлов → 86400 × 50ms = 72 минуты только на открытие файлов!

# ✅ Решение: компакция через Delta/Iceberg
spark.sql("OPTIMIZE delta.`s3a://bucket/events`")
# Объединяет мелкие файлы в крупные (target: 1 GB)

Production-кейс: миграция с FileOutputCommitter на Iceberg

Рассмотрим реальный сценарий: ETL-пайплайн обрабатывает события из Kafka и записывает партиционированные Parquet-файлы в S3. Команда сталкивается с проблемами и мигрирует на Iceberg.

До миграции: симптомы проблем

# Исходный код (legacy):
from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .appName("events-etl-legacy") \
    .config("spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version", "2") \
    .getOrCreate()

def process_batch_legacy(batch_df, batch_id):
    """Исходная обработка - прямая запись Parquet."""

    enriched = batch_df \
        .withColumn("event_date", F.to_date("event_timestamp")) \
        .withColumn("event_hour", F.hour("event_timestamp"))

    # ❌ Проблемы:
    # 1. mode("overwrite") удаляет partitions до записи новых
    # 2. FileOutputCommitter v2 → неатомарный commit
    # 3. При падении → data loss или corruption
    enriched.write \
        .mode("overwrite") \
        .partitionBy("event_date", "event_hour") \
        .parquet("s3a://data-lake/raw/events/")

# Логи показывают:
# 2024-06-01 03:42:15 INFO Committing output of task attempt_0_000_m_000000_0
# 2024-06-01 03:42:15 INFO Committing output of task attempt_0_000_m_000001_0
# ... (повторяется 500 раз)
# 2024-06-01 05:18:33 ERROR Heartbeat timeout for driver
# Job FAILED after 96 minutes

# Анализ: 500 тасков × 2 файла × (CopyObject 100ms + DeleteObject 20ms) = ~2 минуты
# Но с сетевой нестабильностью и throttling → 96 минут

После миграции: Iceberg + правильная конфигурация

from pyspark.sql import SparkSession, functions as F
from pyspark.sql.types import StructType, StructField, StringType, LongType, TimestampType


def create_spark_session() -> SparkSession:
    return SparkSession.builder \
        .appName("events-etl-iceberg") \
        .config("spark.sql.extensions",
                "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
        .config("spark.sql.catalog.glue",
                "org.apache.iceberg.spark.SparkCatalog") \
        .config("spark.sql.catalog.glue.catalog-impl",
                "org.apache.iceberg.aws.glue.GlueCatalog") \
        .config("spark.sql.catalog.glue.warehouse",
                "s3a://data-lake/warehouse/") \
        .config("spark.hadoop.fs.s3a.connection.maximum", "200") \
        .config("spark.hadoop.fs.s3a.multipart.size", "134217728") \
        .getOrCreate()


def ensure_iceberg_table(spark: SparkSession) -> None:
    """DDL для Iceberg таблицы - idempotent."""
    spark.sql("""
        CREATE TABLE IF NOT EXISTS glue.raw.events (
            event_id        STRING  NOT NULL,
            user_id         BIGINT,
            event_type      STRING,
            payload         STRING,
            event_timestamp TIMESTAMP NOT NULL,
            event_date      DATE      NOT NULL,
            event_hour      INT       NOT NULL
        )
        USING iceberg
        PARTITIONED BY (event_date, event_hour)
        TBLPROPERTIES (
            'write.format.default'           = 'parquet',
            'write.parquet.compression-codec'= 'zstd',
            'write.target-file-size-bytes'   = '134217728',
            'format-version'                 = '2'
        )
    """)


def process_batch_iceberg(spark: SparkSession, batch_df, batch_id: int) -> dict:
    """
    Обработка микробатча с записью в Iceberg.

    Что изменилось:
    - Нет _temporary/ директорий
    - Нет rename операций
    - Атомарный commit через snapshot metadata
    - При падении: нет data corruption, orphan files → vacuum удалит
    """
    enriched = batch_df \
        .withColumn("event_date", F.to_date("event_timestamp")) \
        .withColumn("event_hour", F.hour("event_timestamp")) \
        .dropDuplicates(["event_id"])  # идемпотентность при retry

    enriched.createOrReplaceTempView("_batch_events")

    # Iceberg MERGE для exactly-once семантики:
    spark.sql("""
        MERGE INTO glue.raw.events AS target
        USING _batch_events AS source
        ON target.event_id = source.event_id
        WHEN NOT MATCHED THEN INSERT *
    """)

    # Метрики коммита:
    stats = spark.sql(f"""
        SELECT operation, summary
        FROM glue.raw.events.history
        ORDER BY made_current_at DESC
        LIMIT 1
    """).first()

    return {
        "batch_id": batch_id,
        "operation": stats["operation"],
    }


def run_streaming_pipeline() -> None:
    spark = create_spark_session()
    ensure_iceberg_table(spark)

    kafka_df = spark.readStream \
        .format("kafka") \
        .option("kafka.bootstrap.servers", "kafka:9092") \
        .option("subscribe", "events") \
        .load()

    parsed_df = kafka_df.selectExpr(
        "CAST(value AS STRING) AS raw_json",
        "timestamp AS kafka_timestamp"
    ).select(
        F.get_json_object("raw_json", "$.event_id").alias("event_id"),
        F.get_json_object("raw_json", "$.user_id").cast("long").alias("user_id"),
        F.get_json_object("raw_json", "$.event_type").alias("event_type"),
        F.get_json_object("raw_json", "$.payload").alias("payload"),
        F.get_json_object("raw_json", "$.timestamp").cast("timestamp")
            .alias("event_timestamp"),
    )

    query = parsed_df.writeStream \
        .foreachBatch(
            lambda df, bid: process_batch_iceberg(spark, df, bid)
        ) \
        .trigger(processingTime="5 minutes") \
        .option("checkpointLocation", "s3a://data-lake/checkpoints/events/") \
        .start()

    query.awaitTermination()

Сравнение производительности: до и после

Legacy FileOutputCommitter v2 + Parquet:
├── Запись данных (500 параллельных тасков): 8 минут
├── Job Commit (rename 1000 файлов × CopyObject): 96 минут (!)
├── Total pipeline: 104 минуты
├── Успешных запусков из 100: 73 (27% падений из-за timeout)
└── Инциденты data corruption за месяц: 3

Iceberg + Spark:
├── Запись данных (500 параллельных тасков): 8 минут
├── Job Commit (snapshot metadata update): 2 секунды
├── Total pipeline: 8 минут 2 секунды
├── Успешных запусков из 100: 100 (0% падений)
└── Инциденты data corruption за месяц: 0

Ускорение commit-фазы: с 96 минут до 2 секунд - в 2880 раз.

Итого: ключевые выводы

  1. S3 - не файловая система. Это flat key-value store с HTTP API. Слэш в ключе - просто символ, директорий нет, rename - это Copy + Delete.

  2. FileOutputCommitter ломается на S3. v1 - катастрофически медленно (O(N²) операций). v2 - быстрее, но data corruption при падении.

  3. S3A Committers (Magic, Staging) - правильное решение для «сырых» файловых записей. Magic Committer использует Multipart Upload как транзакцию. Staging Committer - локальный диск как буфер.

  4. Lakehouse-форматы (Delta Lake, Iceberg) - радикальное решение: данные пишутся сразу в финальные пути с UUID-именами, коммит - атомарное обновление одного metadata-файла. Rename не нужен вообще.

  5. S3 Throttling - при высокой параллельности Spark может получать 503 Slow Down. Решение: UUID-именование файлов (как в Iceberg/Delta) или salting prefix.

  6. Eventual Consistency - историческая проблема AWS S3, решена в декабре 2020. На self-hosted S3-совместимых хранилищах (MinIO, Ceph) - уточняйте гарантии конкретного продукта.

  7. Практическое правило для production: никогда не использовать стандартный FileOutputCommitter для записи в S3. Либо S3A Committer (для Parquet/CSV без Lakehouse), либо Delta Lake/Iceberg (рекомендуется).