Committer Algorithms: magic committer vs staging committer vs partitioned
Детальный разбор S3A Committers: как Magic, Staging и Partitioned коммитеры работают изнутри, как выбрать правильный алгоритм для вашей нагрузки и как избежать типичных ошибок production.
Commit Protocol - это не деталь реализации¶
Когда разработчики пишут df.write.parquet("s3a://bucket/output/"), они думают о трансформации данных. Commit protocol - механизм, который делает эту запись атомарной - воспринимается как «низкоуровневая деталь Hadoop». На HDFS это действительно так: rename за 50 мс, никаких проблем.
На S3 commit protocol превращается в самую критическую часть пайплайна. Именно он определяет:
- Корректность: увидят ли читатели частичные данные при падении задачи?
- Производительность: сколько времени займёт «коммит» - секунды или часы?
- Стоимость: сколько S3 API-запросов генерирует ваш пайплайн?
- Надёжность: что происходит при retry, speculative execution, OutOfMemory?
Предыдущий урок объяснил, почему стандартный FileOutputCommitter ломается на S3. Этот урок - про что делать: детальный разбор трёх специализированных алгоритмов семейства S3A Committers и практика их применения.
Фундамент: S3 Multipart Upload как примитив атомарности¶
Все три S3A Committer используют один общий механизм - S3 Multipart Upload (MPU). Прежде чем разбирать каждый алгоритм, нужно понять, как работает MPU и почему он даёт то, чего нет у обычного PUT Object.
Обычный PUT Object vs Multipart Upload¶
Обычный PUT Object - синхронная операция: вы отправляете байты, S3 сохраняет объект, он немедленно становится видимым для всех. Если соединение рвётся на середине - объект не создаётся, ничего не сохраняется. Максимальный размер - 5 GB.
Multipart Upload разбивает загрузку на три фазы:
Ключевое свойство MPU: между CreateMultipartUpload и CompleteMultipartUpload объект не существует. Части загружены в S3 (и хранятся, тарифицируются), но недоступны через GET/HEAD/LIST. Только CompleteMultipartUpload «публикует» объект атомарно.
Это и есть механизм атомарности для S3A Committers: экзекуторы загружают данные во время выполнения задачи, а драйвер «публикует» их одной командой при коммите.
# Минимальный размер части MPU: 5 MB (кроме последней части)
# Максимальный размер объекта: 5 TB
# Максимальное количество частей: 10 000
# Настройка размера части в Hadoop S3A:
spark.conf.set("fs.s3a.multipart.size", "134217728") # 128 MB - оптимально
spark.conf.set("fs.s3a.multipart.threshold", "134217728") # начинать MPU от 128 MB
Проблема незавершённых MPU¶
При падении экзекутора или отмене задачи AbortMultipartUpload может не быть вызван. Незавершённые части остаются в S3 и тарифицируются как обычное хранилище. За несколько недель интенсивной работы кластера это может превратиться в сотни гигабайт «мусора».
Обязательная настройка жизненного цикла бакета (через AWS Console или Terraform):
{
"Rules": [
{
"ID": "abort-incomplete-multipart-uploads",
"Status": "Enabled",
"Filter": {
"Prefix": ""
},
"AbortIncompleteMultipartUpload": {
"DaysAfterInitiation": 3
}
}
]
}
# Через boto3:
import boto3
s3 = boto3.client("s3")
s3.put_bucket_lifecycle_configuration(
Bucket="my-data-lake",
LifecycleConfiguration={
"Rules": [{
"ID": "abort-incomplete-mpu",
"Status": "Enabled",
"Filter": {"Prefix": ""},
"AbortIncompleteMultipartUpload": {
"DaysAfterInitiation": 3
}
}]
}
)
Magic Committer: нативная облачная атомарность¶
Идея: использовать MPU как транзакцию¶
Magic Committer - наиболее элегантный из трёх алгоритмов. Его идея: вместо того чтобы имитировать POSIX-семантику (писать во временную директорию, потом переименовывать), использовать нативные примитивы S3 - Multipart Upload - как транзакционный механизм.
Название «Magic» отражает то, как это выглядит со стороны: экзекутор записывает файл, файл «существует» в каком-то промежуточном состоянии, а после job commit - мгновенно появляется в финальном пути. Как будто магия.
Пошаговая механика Magic Committer¶
Разберём каждый шаг подробно:
Шаг 1: CreateMultipartUpload под магическим префиксом.
Экзекутор начинает запись файла не в финальный путь output/part-00042.parquet, а под специальным префиксом __magic/. Этот префикс Hadoop S3A FileSystem перехватывает и превращает в CreateMultipartUpload вместо обычного PUT Object. S3 возвращает UploadId.
Шаг 2: UploadPart. Экзекутор пишет Parquet-данные потоком, нарезая на части по 128 MB. Каждая часть уходит в S3 через UploadPart. Данные физически в S3, но объект не существует - только незавершённые части.
Шаг 3: Сохранение Pending Commit. Вместо вызова CompleteMultipartUpload (которое бы «опубликовало» объект), экзекутор сохраняет крошечный JSON-файл - «инструкцию коммита»: финальный путь + UploadId + список ETag всех частей. Этот JSON пишется в отдельную директорию.
Шаг 4: Job Commit на драйвере.
После успешного завершения всех задач драйвер собирает все JSON-инструкции и вызывает CompleteMultipartUpload для каждой. Каждый вызов - HTTP-запрос к S3 API, который занимает ~50–200 мс. Никакого копирования данных! Объекты «появляются» атомарно в финальных путях.
Особенность: speculative execution¶
Если Spark запускает спекулятивную копию задачи (на другом экзекуторе), оба экзекутора создают разные MPU с разными UploadId. При task commit коммитер использует механизм «last writer wins»: сохраняется JSON только для победившего таска. При job commit вызывается CompleteMultipartUpload только для победителя. Для проигравшего - AbortMultipartUpload (части удаляются).
Это корректная семантика exactly-once: в финальный путь попадают данные ровно одной задачи.
Конфигурация Magic Committer¶
from pyspark.sql import SparkSession
def create_spark_magic_committer(
app_name: str,
s3_endpoint: str = None
) -> SparkSession:
"""
SparkSession с Magic Committer.
Требования:
- Hadoop >= 3.1.1 (модуль hadoop-aws)
- S3 с поддержкой Multipart Upload (AWS S3, MinIO >= 2022)
"""
builder = SparkSession.builder.appName(app_name)
# Обязательные настройки Magic Committer:
builder = builder \
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.fs.s3a.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")
# Для всех бакетов:
builder = builder \
.config("spark.hadoop.fs.s3a.bucket.all.committer.magic.enabled", "true")
# Производительность MPU:
builder = builder \
.config("spark.hadoop.fs.s3a.multipart.size", "134217728") \ # 128 MB
.config("spark.hadoop.fs.s3a.multipart.threshold", "134217728") \ # порог MPU
.config("spark.hadoop.fs.s3a.threads.max", "64") \
.config("spark.hadoop.fs.s3a.connection.maximum", "200")
if s3_endpoint:
# MinIO / self-hosted S3-compatible:
builder = builder \
.config("spark.hadoop.fs.s3a.endpoint", s3_endpoint) \
.config("spark.hadoop.fs.s3a.path.style.access", "true")
return builder.getOrCreate()
# Пример использования:
spark = create_spark_magic_committer("etl-magic", "http://minio:9000")
df = spark.read.parquet("s3a://raw/events/")
df.write \
.mode("append") \
.partitionBy("event_date") \
.parquet("s3a://processed/events/")
# → Нет rename, commit = N × CompleteMultipartUpload HTTP запросов
Плюсы и минусы Magic Committer¶
Плюсы:
- Самый быстрый commit из трёх (только HTTP метаданные, нет IO данных).
- Не требует дополнительной инфраструктуры (нет нужды в HDFS, ZooKeeper).
- Идеален для Spark на Kubernetes, AWS EMR, Databricks (serverless без локальных дисков).
- Масштабируется линейно: 1000 файлов = 1000 CompleteMultipartUpload запросов параллельно.
Минусы:
- Требует поддержки MPU в S3-хранилище. Старые версии MinIO/Ceph могут иметь баги.
- При падении оставляет незавершённые MPU (нужны lifecycle rules).
- Не работает с хранилищами, которые не поддерживают Multipart Upload (некоторые совместимые API, например ранние версии OpenStack Swift).
- Сложнее отлаживать: «магические» пути в
__magic/сбивают с толку при ручном просмотре бакета.
Staging Committer: надёжность через локальный буфер¶
Идея: сначала POSIX, потом S3¶
Staging Committer (также известный как Directory Committer) использует другой подход: вместо того чтобы искать атомарность в примитивах S3, он использует локальную POSIX-файловую систему диска экзекутора как промежуточный буфер. Все данные сначала пишутся локально (с настоящим атомарным rename), и только при коммите загружаются в S3.
Этот подход разработала команда Netflix для своей Hadoop-платформы. Netflix оперирует петабайтами данных в S3 и накопила огромный опыт с edge cases commit protocols.
Пошаговая механика Staging Committer¶
Шаг 1: Запись на локальный диск.
Экзекутор пишет Parquet/ORC данные не в S3, а в spark.hadoop.fs.s3a.committer.staging.tmp.path (обычно /tmp/spark-staging/<job_id>/task_<id>/). Это обычная запись в локальную файловую систему - POSIX, быстро, надёжно.
Шаг 2: Task Commit на локальной POSIX FS.
После успешного завершения вычислений экзекутор делает rename в локальной файловой системе: task_tmp_path → task_committed_path. На POSIX это атомарная операция, занимает микросекунды. Если задача упала - мусор остаётся в task_tmp_path, который при повторном запуске перезаписывается.
Шаг 3: Сохранение манифеста. Экзекутор записывает манифест-файл в координирующее хранилище (HDFS или общую директорию): «задача 42 записала файл по такому-то локальному пути». Это крошечный JSON.
Шаг 4: Job Commit - параллельная загрузка в S3.
Драйвер собирает все манифесты. Для каждого задача - загружает данные с локального диска экзекутора напрямую в S3 через PUT Object или Multipart Upload, сразу в финальный путь. Данные идут прямо из /tmp/ в s3a://bucket/output/part-XXXXX.parquet.
Конфигурация Staging Committer¶
def create_spark_staging_committer(
app_name: str,
staging_dir: str = "/tmp/spark-staging",
conflict_mode: str = "replace",
s3_endpoint: str = None
) -> SparkSession:
"""
SparkSession с Directory Staging Committer.
conflict_mode:
- 'replace' : перезаписать партиции из текущего batch (безопасно)
- 'append' : добавить файлы к существующим
- 'fail' : ошибка если партиция уже существует
"""
builder = SparkSession.builder.appName(app_name)
builder = builder \
.config("spark.hadoop.fs.s3a.committer.name", "directory") \
.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") \
.config("spark.hadoop.fs.s3a.committer.staging.tmp.path", staging_dir) \
.config("spark.hadoop.fs.s3a.committer.staging.conflict-mode",
conflict_mode) \
.config("spark.hadoop.fs.s3a.committer.staging.unique-filenames", "true")
# Производительность загрузки в S3 при job commit:
builder = builder \
.config("spark.hadoop.fs.s3a.threads.max", "64") \
.config("spark.hadoop.fs.s3a.connection.maximum", "200") \
.config("spark.hadoop.fs.s3a.multipart.size", "134217728") \
.config("spark.hadoop.fs.s3a.fast.upload", "true") \
.config("spark.hadoop.fs.s3a.fast.upload.buffer", "bytebuffer")
if s3_endpoint:
builder = builder \
.config("spark.hadoop.fs.s3a.endpoint", s3_endpoint) \
.config("spark.hadoop.fs.s3a.path.style.access", "true")
return builder.getOrCreate()
Управление локальным дисковым пространством¶
Staging Committer требует локального дискового пространства, равного размеру output batch. Это критическое ограничение для больших датасетов.
import os
import shutil
from pyspark.sql import SparkSession
def calculate_staging_space_gb(
df,
compression_ratio: float = 0.3 # Parquet zstd обычно 3x сжатие
) -> float:
"""
Оценить необходимое место для staging.
Приблизительно: размер данных × (1 - compression_ratio).
"""
# Spark не знает размер до записи, оцениваем по схеме и count():
num_rows = df.count()
num_cols = len(df.columns)
estimated_bytes = num_rows * num_cols * 8 # rough estimate
estimated_parquet_bytes = estimated_bytes * compression_ratio
return estimated_parquet_bytes / (1024 ** 3)
def check_staging_space(staging_dir: str, required_gb: float) -> bool:
"""Проверить достаточность места перед запуском Staging Committer."""
stat = shutil.disk_usage(staging_dir)
available_gb = stat.free / (1024 ** 3)
if available_gb < required_gb * 1.2: # 20% запас
raise RuntimeError(
f"Недостаточно места для Staging Committer: "
f"нужно {required_gb:.1f} GB, доступно {available_gb:.1f} GB. "
f"Рассмотрите Magic Committer или уменьшите размер batch."
)
return True
def cleanup_staging(staging_dir: str, job_id: str) -> None:
"""Очистить staging директорию после коммита."""
job_dir = os.path.join(staging_dir, job_id)
if os.path.exists(job_dir):
shutil.rmtree(job_dir)
Плюсы и минусы Staging Committer¶
Плюсы:
- Максимальная надёжность: при падении задачи в S3 ничего не попадает (мусор только на локальном диске, который очищается при повторном запуске).
- Детерминированный commit: нет race conditions, нет проблем с MPU мусором.
- Лучшая совместимость: работает с любым S3-совместимым хранилищем, даже без поддержки MPU.
- Понятная отладка: можно руками посмотреть файлы в
/tmp/до и после коммита.
Минусы:
- Требует локальное дисковое пространство ≈ размеру batch (критично для больших датасетов).
- Double write: данные пишутся дважды (локально + в S3), двойная нагрузка на сеть экзекуторов.
- Не подходит для Spark на Kubernetes с ephemeral pods без persistent volumes.
- Job commit - sequential по умолчанию: загрузка в S3 происходит при commit, не параллельно с вычислениями.
Partitioned Committer: точечная перезапись партиций¶
Проблема overwrite партиционированных таблиц¶
Самая частая операция в batch ETL-пайплайнах - перезапись конкретных партиций. Например, ежедневный джоб пересчитывает данные за вчера и перезаписывает партицию date=2024-06-15.
Стандартный FileOutputCommitter при mode("overwrite") делает следующее:
- Удаляет всю выходную директорию (
DeleteObjectsдля всех ключей с префиксомoutput/). - Записывает новые данные.
Это опасно: если джоб пишет только партицию date=2024-06-15, он всё равно удалит данные за все остальные даты. Это классический «случайный DROP TABLE».
Spark добавил spark.sql.sources.partitionOverwriteMode = dynamic для решения этой проблемы - режим, при котором перезаписываются только партиции, в которые реально пишут экзекуторы. Но реализация через FileOutputCommitter на S3 всё равно использует rename.
Partitioned Committer - это подвид Staging Committer, специализированный для динамической перезаписи партиций с корректной семантикой на S3.
Механика Partitioned Committer¶
Partitioned Committer отслеживает, в какие именно партиции пишет каждый экзекутор, и при job commit выполняет точечную замену: удаляет старые файлы только в затронутых партициях, загружает новые. Незатронутые партиции не трогаются.
Конфигурация Partitioned Committer¶
def create_spark_partitioned_committer(
app_name: str,
staging_dir: str = "/tmp/spark-staging",
s3_endpoint: str = None
) -> SparkSession:
"""
SparkSession с Partitioned Staging Committer.
Специально для partitionBy + mode('overwrite') на S3.
"""
builder = SparkSession.builder.appName(app_name)
builder = builder \
.config("spark.hadoop.fs.s3a.committer.name", "partitioned") \
.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") \
.config("spark.hadoop.fs.s3a.committer.staging.tmp.path", staging_dir) \
.config("spark.hadoop.fs.s3a.committer.staging.conflict-mode", "replace") \
.config("spark.sql.sources.partitionOverwriteMode", "dynamic")
if s3_endpoint:
builder = builder \
.config("spark.hadoop.fs.s3a.endpoint", s3_endpoint) \
.config("spark.hadoop.fs.s3a.path.style.access", "true")
return builder.getOrCreate()
# Пример: перезапись конкретных партиций
def overwrite_partitions(
spark: SparkSession,
source_df,
output_path: str,
partition_cols: list
) -> None:
"""
Безопасная перезапись партиций на S3 через Partitioned Committer.
Только затронутые партиции будут перезаписаны.
"""
(
source_df
.write
.mode("overwrite")
.partitionBy(*partition_cols)
.parquet(output_path)
)
# Partitioned Committer гарантирует:
# 1. Только партиции из source_df будут перезаписаны
# 2. Остальные партиции не тронуты
# 3. Если джоб упадёт - никакие данные в S3 не будут повреждены
Режимы conflict-mode¶
Параметр spark.hadoop.fs.s3a.committer.staging.conflict-mode определяет поведение при конфликте между existing и incoming данными:
# conflict-mode = 'replace' (рекомендуется для overwrite):
# Удалить существующие файлы в партиции → загрузить новые
# Использовать для: ежедневные пересчёты, исправления данных
# conflict-mode = 'append' (для incremental loads):
# Добавить новые файлы к существующим в партиции
# Использовать для: добавление новых данных без перезаписи старых
# conflict-mode = 'fail' (для защиты от случайного overwrite):
# Ошибка если партиция уже содержит данные
# Использовать для: first-write pipelines, защита critical data
# Динамическое управление через код:
def write_with_conflict_mode(df, output_path, partition_cols, mode):
df.write \
.mode("overwrite") \
.option("partitionOverwriteMode", "dynamic") \
.partitionBy(*partition_cols) \
.parquet(output_path)
Коммитеры и Lakehouse-форматы: как они взаимодействуют¶
Студенты часто задают вопрос: «Если я использую Delta Lake или Iceberg, нужны ли мне S3A Committers?»
Ответ неочевидный: Lakehouse-форматы имеют собственные commit protocols, которые по умолчанию не используют S3A Committers. Но это не значит, что знание о коммитерах не нужно.
Когда S3A Committers нужны даже с Lakehouse¶
Временные промежуточные файлы: даже в Iceberg-пайплайне Spark может писать промежуточные Parquet-файлы (например, при sort-merge join с материализацией результата). Эти промежуточные файлы идут через FileOutputCommitter.
Spark Structured Streaming checkpoint: checkpointing в стриминг пишет через стандартный FileSystem API - без S3A Committer это медленно.
Внешние инструменты: если рядом с Lakehouse работают инструменты, которые пишут через Spark FileSystem API напрямую (legacy ETL, AWS Glue с Parquet), им нужны S3A Committers.
Вывод: при работе с Delta Lake / Iceberg S3A Committers не нужны для самих таблиц, но рекомендуются как базовая настройка кластера.
Anti-patterns: как сломать pipeline неправильным коммитером¶
Антипаттерн 1: FileOutputCommitter v1 на больших датасетах¶
# ❌ Дефолтная конфигурация Spark на S3:
spark = SparkSession.builder \
.appName("legacy-etl") \
# Нет конфигурации коммитера → используется FileOutputCommitter v1
.getOrCreate()
df.write \
.mode("overwrite") \
.partitionBy("date") \
.parquet("s3a://bucket/output/")
# Что происходит:
# 1. Запись 500 файлов в _temporary/attempt_XXX/ → OK
# 2. Task commits: rename attempt → task (×500 файлов в S3) → N×CopyObject
# 3. Job commit: rename task → output (×500 файлов в S3) → ещё N×CopyObject
# 4. Итого: 1000 CopyObject операций + 1000 DeleteObject
# 5. Время коммита: 10–60 минут при 500 файлах по 500MB
# 6. Вероятность timeout на драйвере: высокая
Антипаттерн 2: Magic Committer на хранилищах без надёжного MPU¶
# ❌ Magic Committer на старом MinIO с багами MPU:
spark.conf.set("spark.hadoop.fs.s3a.committer.name", "magic")
# Проблема на старых версиях MinIO (до 2022):
# CompleteMultipartUpload может вернуть 200 OK, но объект создаётся с задержкой.
# Spark думает, что коммит прошёл успешно,
# но читатели ещё не видят файлы → пустой результат.
# Диагностика:
# 2024-06-15 03:42:00 INFO MagicCommitTracker: completing multipart upload...
# 2024-06-15 03:42:01 INFO MagicCommitTracker: upload complete
# 2024-06-15 03:42:01 INFO Job committed successfully
# Но: SELECT COUNT(*) FROM table → 0 строк!
# ✅ Решение: использовать Directory Staging Committer для on-premise хранилищ
spark.conf.set("spark.hadoop.fs.s3a.committer.name", "directory")
Антипаттерн 3: Staging Committer без проверки дискового места¶
# ❌ Staging Committer без мониторинга локальных дисков:
spark.conf.set("spark.hadoop.fs.s3a.committer.name", "directory")
spark.conf.set("spark.hadoop.fs.s3a.committer.staging.tmp.path", "/tmp/")
# Джоб обрабатывает 2TB данных на 20 воркерах = 100 GB на каждый воркер
# Если диск воркера = 50 GB → No space left on device
# Задача падает с невнятной ошибкой:
# java.io.IOException: No space left on device
# Все данные на локальных дисках → 0 байт в S3
# ✅ Решение: использовать NVMe instance storage (большие диски),
# или Magic Committer, или уменьшить размер batch
spark.conf.set("spark.hadoop.fs.s3a.committer.staging.tmp.path",
"/mnt/nvme/spark-staging/") # NVMe диск 1-2 TB
Антипаттерн 4: Partitioned Committer с conflict-mode='fail' для overwrite¶
# ❌ Неправильный conflict-mode для ежедневного пересчёта:
spark.conf.set("spark.hadoop.fs.s3a.committer.name", "partitioned")
spark.conf.set("spark.hadoop.fs.s3a.committer.staging.conflict-mode", "fail")
# Ежедневный джоб пересчитывает вчерашние данные:
df_yesterday.write \
.mode("overwrite") \
.partitionBy("date") \
.parquet("s3a://bucket/events/")
# При первом запуске: OK (партиции не существуют)
# При повторном запуске (retry, исправление данных):
# org.apache.hadoop.fs.s3a.commit.CommitConflictException:
# Partition date=2024-06-15 already exists and conflict mode is 'fail'
# ✅ Для перезаписи: conflict-mode='replace'
spark.conf.set("spark.hadoop.fs.s3a.committer.staging.conflict-mode", "replace")
Антипаттерн 5: Игнорирование незавершённых MPU¶
# ❌ Magic Committer без lifecycle rules на бакете:
# После интенсивной недели разработки (много failed jobs):
# - Незавершённые MPU: 500 GB
# - Ежемесячный счёт S3: +$11.50 (S3 Standard = $0.023 / GB)
# Через 3 месяца: 500 GB × 3 = 1500 GB лишних расходов
# Проверка незавершённых MPU:
import boto3
s3 = boto3.client("s3")
response = s3.list_multipart_uploads(Bucket="my-data-lake")
if "Uploads" in response:
for upload in response["Uploads"]:
print(
f"Key: {upload['Key']}, "
f"UploadId: {upload['UploadId']}, "
f"Initiated: {upload['Initiated']}"
)
print(f"Всего незавершённых MPU: {len(response['Uploads'])}")
def abort_old_multipart_uploads(bucket: str, days_old: int = 3) -> int:
"""Принудительно отменить все MPU старше N дней."""
from datetime import datetime, timezone, timedelta
s3 = boto3.client("s3")
cutoff = datetime.now(timezone.utc) - timedelta(days=days_old)
aborted = 0
paginator = s3.get_paginator("list_multipart_uploads")
for page in paginator.paginate(Bucket=bucket):
for upload in page.get("Uploads", []):
if upload["Initiated"] < cutoff:
s3.abort_multipart_upload(
Bucket=bucket,
Key=upload["Key"],
UploadId=upload["UploadId"]
)
aborted += 1
return aborted
Диагностика: как увидеть работу коммитера¶
Spark UI: признаки правильного коммитера¶
В Spark UI (вкладка Stages) можно увидеть, как долго длился финальный job commit:
Stage 5 (коммит):
Duration: 00:02:15 ← ❌ FileOutputCommitter v1: 2 минуты на rename
Duration: 00:00:02 ← ✅ Magic Committer: 2 секунды на CompleteMultipartUpload
Duration: 00:00:45 ← ⚠️ Staging Committer: 45 секунд на PUT в S3
Логи драйвера: активность коммитера¶
# Включить verbose логирование S3A Committer:
import logging
logging.getLogger("org.apache.hadoop.fs.s3a.commit").setLevel(logging.DEBUG)
# Magic Committer в логах:
# INFO MagicCommitTracker: Starting commit of 234 pending uploads
# INFO MagicCommitTracker: Completing upload for output/part-00000.parquet
# INFO MagicCommitTracker: Completed 234/234 uploads in 1.8s
# INFO _MagicOutputCommitter: Job committed in 1.8s
# Staging Committer в логах:
# INFO StagingCommitter: Uploading 234 files from local staging to S3
# INFO StagingCommitter: PUT output/part-00000.parquet (234MB)
# INFO StagingCommitter: Uploaded 234/234 files in 42.3s
Мониторинг S3 операций через CloudWatch¶
# AWS CloudWatch метрики для диагностики:
import boto3
from datetime import datetime, timedelta
cloudwatch = boto3.client("cloudwatch")
# Количество CopyObject запросов (признак legacy FileOutputCommitter):
response = cloudwatch.get_metric_statistics(
Namespace="AWS/S3",
MetricName="NumberOfObjects",
Dimensions=[{"Name": "BucketName", "Value": "my-data-lake"}],
StartTime=datetime.now() - timedelta(hours=1),
EndTime=datetime.now(),
Period=300,
Statistics=["Sum"]
)
# Ключевые метрики для мониторинга:
# - PutRequests: количество PUT (ожидаем рост при job commit)
# - CopyRequests: количество COPY (должно быть 0 при Magic/Staging!)
# - DeleteRequests: количество DELETE при overwrite
Production-кейс: выбор коммитера для разных нагрузок¶
Сценарий A: Spark на Kubernetes без persistent storage¶
# Оптимальная конфигурация для Spark на Kubernetes:
spark = SparkSession.builder \
.appName("k8s-etl") \
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.fs.s3a.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") \
.config("spark.hadoop.fs.s3a.bucket.all.committer.magic.enabled", "true") \
.getOrCreate()
Сценарий B: Spark на bare-metal с NVMe-дисками¶
Кластер с физическими серверами, каждый экзекутор имеет 2 TB NVMe-диска. Задача: ежедневный пересчёт партиций за прошедшую неделю.
# Оптимальная конфигурация для bare-metal с NVMe:
spark = SparkSession.builder \
.appName("bare-metal-daily-recalc") \
.config("spark.hadoop.fs.s3a.committer.name", "partitioned") \
.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") \
.config("spark.hadoop.fs.s3a.committer.staging.tmp.path",
"/mnt/nvme/spark-staging") \
.config("spark.hadoop.fs.s3a.committer.staging.conflict-mode", "replace") \
.config("spark.sql.sources.partitionOverwriteMode", "dynamic") \
.getOrCreate()
# Перезапись партиций за неделю:
def daily_recalc(spark: SparkSession, start_date: str, end_date: str) -> None:
raw = spark.read.parquet(f"s3a://raw/events/")
processed = raw \
.filter(f"event_date BETWEEN '{start_date}' AND '{end_date}'") \
.groupBy("event_date", "user_id") \
.agg({"amount": "sum", "event_id": "count"})
# Partitioned Committer перезапишет только даты из start_date..end_date,
# не трогая исторические партиции
processed.write \
.mode("overwrite") \
.partitionBy("event_date") \
.parquet("s3a://processed/daily-stats/")
Сценарий C: AWS EMR с динамическими кластерами¶
# EMR bootstrap action / конфигурация кластера:
# В emr-config.json:
# {
# "Classification": "spark-defaults",
# "Properties": {
# "spark.hadoop.fs.s3a.committer.name": "magic",
# "spark.hadoop.fs.s3a.committer.magic.enabled": "true",
# ...
# }
# }
# В коде пайплайна:
def create_emr_spark_session(app_name: str) -> SparkSession:
"""
На AWS EMR конфигурация обычно задаётся на уровне кластера.
Здесь - минимальные переопределения.
"""
return SparkSession.builder \
.appName(app_name) \
.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()
Сводная матрица выбора коммитера¶
def recommend_committer(
environment: str, # 'kubernetes' | 'emr' | 'bare-metal' | 'local'
local_disk_gb: int, # доступное место на экзекуторе
batch_size_gb: int, # ожидаемый размер output
s3_type: str, # 'aws' | 'minio-new' | 'minio-old' | 'ceph'
operation: str, # 'append' | 'overwrite-partition' | 'overwrite-full'
use_lakehouse: bool # Delta Lake / Iceberg
) -> dict:
"""
Рекомендатор коммитера на основе параметров среды.
"""
if use_lakehouse:
return {
"committer": "встроенный в Lakehouse-формат",
"notes": "Delta/Iceberg имеют собственные commit protocols. "
"S3A Committer нужен только для промежуточных файлов.",
"config": {"spark.hadoop.fs.s3a.committer.name": "magic"}
}
if environment == "kubernetes" and local_disk_gb < batch_size_gb:
if s3_type in ("aws", "minio-new"):
return {
"committer": "magic",
"reason": "Нет места для Staging, MPU надёжен",
"config": {
"spark.hadoop.fs.s3a.committer.name": "magic",
"spark.hadoop.fs.s3a.committer.magic.enabled": "true"
}
}
else:
return {
"committer": "iceberg/delta",
"reason": "Нет места для Staging, MPU ненадёжен → Lakehouse",
"config": {}
}
if operation == "overwrite-partition" and local_disk_gb >= batch_size_gb:
return {
"committer": "partitioned",
"reason": "Оптимален для точечной перезаписи партиций",
"config": {
"spark.hadoop.fs.s3a.committer.name": "partitioned",
"spark.hadoop.fs.s3a.committer.staging.conflict-mode": "replace",
"spark.sql.sources.partitionOverwriteMode": "dynamic"
}
}
if local_disk_gb >= batch_size_gb * 1.5 and s3_type in ("minio-old", "ceph"):
return {
"committer": "directory",
"reason": "Staging надёжнее на хранилищах с нестабильным MPU",
"config": {
"spark.hadoop.fs.s3a.committer.name": "directory",
"spark.hadoop.fs.s3a.committer.staging.conflict-mode": "replace"
}
}
return {
"committer": "magic",
"reason": "Универсальный выбор для большинства случаев",
"config": {
"spark.hadoop.fs.s3a.committer.name": "magic",
"spark.hadoop.fs.s3a.committer.magic.enabled": "true"
}
}
Полная конфигурация production-кластера¶
Ниже - эталонная конфигурация для production Spark-кластера с S3 (или MinIO). Применяется как базовая для всех пайплайнов кластера:
SPARK_S3_PRODUCTION_CONFIG = {
# === Committer (выбрать один) ===
"spark.hadoop.fs.s3a.committer.name": "magic",
"spark.hadoop.fs.s3a.committer.magic.enabled": "true",
"spark.hadoop.fs.s3a.bucket.all.committer.magic.enabled": "true",
"spark.sql.sources.commitProtocolClass":
"org.apache.spark.internal.io.cloud.PathOutputCommitProtocol",
"spark.sql.parquet.output.committer.class":
"org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter",
# === Производительность I/O ===
"spark.hadoop.fs.s3a.connection.maximum": "200",
"spark.hadoop.fs.s3a.threads.max": "64",
"spark.hadoop.fs.s3a.multipart.size": "134217728", # 128 MB
"spark.hadoop.fs.s3a.multipart.threshold": "134217728", # 128 MB
"spark.hadoop.fs.s3a.fast.upload": "true",
"spark.hadoop.fs.s3a.fast.upload.buffer": "array",
"spark.hadoop.fs.s3a.block.size": "134217728", # 128 MB
"spark.hadoop.fs.s3a.readahead.range": "2097152", # 2 MB
# === Retry и устойчивость ===
"spark.hadoop.fs.s3a.retry.limit": "7",
"spark.hadoop.fs.s3a.retry.interval": "500ms",
"spark.hadoop.fs.s3a.attempts.maximum": "20",
"spark.hadoop.fs.s3a.connection.timeout": "200000",
"spark.hadoop.fs.s3a.socket.recv.buffer": "65536",
"spark.hadoop.fs.s3a.connection.establish.timeout": "15000",
# === Отключить legacy FileOutputCommitter ===
"spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version": "2",
"spark.hadoop.mapreduce.fileoutputcommitter.cleanup-failures.ignored": "true",
# === Оптимизация листинга ===
"spark.hadoop.fs.s3a.list.version": "2", # ListObjectsV2
"spark.hadoop.fs.s3a.path.style.access": "false", # true для MinIO
"spark.hadoop.fs.s3a.directory.marker.retention": "keep", # Iceberg/Delta
# === Метрики ===
"spark.hadoop.fs.s3a.experimental.input.fadvise": "sequential", # потоковое чтение
"spark.hadoop.fs.s3a.metrics.enabled": "true",
}
def apply_s3_config(spark: SparkSession, overrides: dict = None) -> SparkSession:
"""Применить production S3 конфигурацию к существующей сессии."""
config = {**SPARK_S3_PRODUCTION_CONFIG, **(overrides or {})}
for key, value in config.items():
spark.conf.set(key, str(value))
return spark
Итого: как выбрать правильный коммитер¶
| Критерий | Magic Committer | Directory Staging | Partitioned Staging |
|---|---|---|---|
| Скорость commit | ✅ Секунды | ⚠️ Минуты (PUT в S3) | ⚠️ Минуты |
| Атомарность | ✅ MPU | ✅ POSIX rename | ✅ POSIX rename |
| Локальный диск | ✅ Не нужен | ❌ Нужен (≥ batch size) | ❌ Нужен |
| Kubernetes | ✅ Идеально | ❌ Нет места | ❌ Нет места |
| Bare-metal NVMe | ✅ | ✅ | ✅ |
| AWS EMR / Databricks | ✅ | ✅ | ✅ |
| On-premise MinIO (старый) | ⚠️ Может глючить | ✅ Надёжно | ✅ Надёжно |
| Overwrite partition | ✅ | ⚠️ conflict-mode=replace | ✅ Специально для этого |
| Overwrite full table | ✅ | ✅ | ⚠️ Избыточен |
| При падении → мусор в S3 | ⚠️ Незавершённые MPU | ✅ Только на локальном диске | ✅ Только на локальном диске |
Практическое правило:
- Kubernetes / serverless → Magic Committer (нет локальных дисков).
- Bare-metal + ежедневные пересчёты партиций → Partitioned Committer.
- Bare-metal + MinIO (старый) + максимальная надёжность → Directory Staging.
- Delta Lake / Apache Iceberg → встроенный commit protocol; Magic Committer как базовая настройка кластера для промежуточных файлов.
- Legacy pipeline, не можете менять конфигурацию → хотя бы v2 FileOutputCommitter + lifecycle rules для cleanup.