S3A Connector: архитектура, конфигурация и параметры производительности
Глубокий разбор S3A Connector: как Spark работает с S3-compatible object storage через Hadoop FileSystem API, почему rename ломает Spark на S3, как работают S3A committers и как настроить production-ready пайплайн для MinIO, Ceph и AWS S3.
Архитектурный сдвиг: от HDFS к Object Storage¶
Ещё десять лет назад стандартной платформой для хранения данных в Hadoop-кластерах был HDFS - распределённая файловая система, физически размещённая на тех же серверах, что и вычислительные ноды. Данные жили рядом с кодом, и Spark мог читать блоки с нулевой или минимальной сетевой передачей (data locality).
Сегодня индустрия сделала кардинальный разворот в сторону object storage: AWS S3, Google Cloud Storage, Azure ADLS Gen2, а в российской инфраструктуре - MinIO, Ceph RGW, Cloud.ru Object Storage, Yandex Object Storage, Selectel и VK Cloud. Причины этого перехода - экономические и архитектурные:
1. Разделение compute и storage (Decoupled Architecture). HDFS требует, чтобы диски жили на тех же серверах, что и вычислительные ноды. Нужно больше вычислений - добавляешь серверы с дисками. Нужно больше хранения - добавляешь серверы с дисками и CPU. Это неэффективно: CPU простаивает, когда нагрузка невысока, но вы платите за всё железо. Object Storage хранит данные независимо от вычислительного кластера. Spark-кластер можно запустить на 10 минут, обработать петабайт данных и выключить - данные остаются в хранилище.
2. Масштабирование без ограничений NameNode. HDFS ограничен объёмом RAM NameNode (каждый файл/блок ≈ 150 байт в памяти). Object Storage масштабируется до экзабайт: у систем вроде MinIO в distributed-режиме метаданные сами по себе распределены по нодам.
3. Stateless вычислительные кластеры. Spark-кластер на Kubernetes, использующий MinIO или S3 как источник и приёмник данных, не имеет локального состояния. Его можно пересоздать за секунды, обновить версию Spark, откатить - данные никуда не денутся.
Но object storage - это не файловая система. И это ключевое понимание, без которого настройка Spark на S3 превращается в бесконечный debugging.
Object Storage как key-value хранилище¶
S3 и его аналоги - это не файловая система POSIX. Это хранилище ключей и значений (key-value store), где ключ - строка вида bucket/prefix/file.parquet, а значение - бинарный blob. Никакого дерева директорий не существует. Когда вы видите «папки» в MinIO UI или AWS Console - это иллюзия, созданная через фильтрацию по префиксу.
Следствия из этой природы:
- Нет atomic rename. В POSIX rename - это операция O(1) над inode-записью в памяти файловой системы. В S3 rename объекта - это
Copy(src, dst)+Delete(src). Для директории с 10 000 файлов это 20 000 HTTP-запросов. При параллельной записи из сотни Executor'ов это убивает производительность. - LIST операции дорогие.
ls /bucket/table/запускает paginated LIST-запрос, который возвращает до 1000 объектов за запрос. При миллионе файлов - 1000 HTTP-запросов только чтобы узнать содержимое директории. - Метадатаческие операции имеют latency. Даже простой
HEADзапрос (проверить существование объекта) занимает 5–50 ms. В HDFS аналогичная операция - 1–2 ms через NameNode в памяти. - Параллельные записи не видели друг друга мгновенно (в старых версиях S3 - eventually consistent LIST). С декабря 2020 года AWS S3 обеспечивает strong consistency, но не все S3-совместимые системы это гарантируют - особенно старые версии Ceph.
Схема выше иллюстрирует ключевое архитектурное различие: в HDFS rename - это мгновенная операция в памяти NameNode, в object storage - это физическое копирование всех объектов. Это различие определяет весь дизайн write path в Spark на S3.
Как устроен S3 изнутри¶
Чтобы понять, почему Spark ведёт себя именно так при работе с S3, нужно разобраться в трёх вещах: как S3 организует данные, как выглядят реальные HTTP-запросы к S3, и как S3 физически хранит объекты внутри AWS.
Data Model: бакеты, объекты и ключи¶
Бакет (Bucket) - это контейнер для объектов, привязанный к конкретному региону AWS (или ноде MinIO). Бакет имеет уникальное имя в глобальном пространстве имён (на AWS - уникальное по всему миру, на MinIO - в пределах кластера). Бакет хранит метаданные политик доступа, настройки версионирования и lifecycle rules.
Объект (Object) - это пара (key → value + metadata):
- Key - строка произвольной длины, например
warehouse/events/year=2024/month=01/part-00001-abc.parquet. Это не путь в файловой системе - это просто строка. Слэши в ключе не создают директорий: они просто часть имени. «Папка»warehouse/events/- это фильтрListObjectsV2с параметромprefix=warehouse/events/иdelimiter=/. - Value - бинарный blob данных, от 0 байт до 5 TB (через Multipart Upload).
- Metadata - системные атрибуты (
Content-Type,Content-Length,ETag,Last-Modified) и пользовательские парыx-amz-meta-*. - ETag - MD5 (или составной MD5 для Multipart Upload) от содержимого объекта. Используется для проверки целостности при GET.
Иллюзия директорий. Когда вы читаете s3a://bucket/warehouse/events/, S3A Connector вызывает ListObjectsV2(prefix=warehouse/events/, delimiter=/). S3 возвращает два типа записей: CommonPrefixes (то, что выглядит как «папки») и Contents (реальные объекты с данным префиксом). Именно поэтому ls на S3 с миллионом объектов - серия HTTP-запросов по 1000 объектов за раз.
S3 REST API: реальные HTTP-запросы¶
S3 - это HTTP REST API. Каждая операция - один или несколько HTTP-запросов. Понимание их структуры объясняет, почему одни операции дешёвые, а другие - дорогие.
PUT Object - загрузка одного объекта (до 5 GB):
PUT /warehouse/events/part-001.parquet HTTP/1.1
Host: bucket.s3.eu-central-1.amazonaws.com
Content-Length: 134217728
Content-Type: application/octet-stream
x-amz-content-sha256: e3b0c44298fc1c149afbf4c8996fb924...
Authorization: AWS4-HMAC-SHA256 Credential=AKID/20240116/eu-central-1/s3/aws4_request, ...
[134 MB binary data]
Ответ: 200 OK с ETag заголовком - MD5 загруженного объекта.
GET Object - скачивание объекта (или диапазона байт):
GET /warehouse/events/part-001.parquet HTTP/1.1
Host: bucket.s3.eu-central-1.amazonaws.com
Range: bytes=0-65535
Authorization: AWS4-HMAC-SHA256 ...
Range - критически важный заголовок. S3A Connector использует range requests для читения только нужных байт объекта. Именно так работает Column Pruning в Parquet: Spark читает footer файла (последние N байт), из него узнаёт смещения нужных column chunks, затем делает range GET только для этих байт. Без поддержки range requests эффективный column pruning был бы невозможен.
ListObjectsV2 - листинг объектов по префиксу:
GET /?list-type=2&prefix=warehouse%2Fevents%2F&delimiter=%2F&max-keys=1000 HTTP/1.1
Host: bucket.s3.eu-central-1.amazonaws.com
Authorization: AWS4-HMAC-SHA256 ...
Ответ содержит до 1000 объектов за раз + NextContinuationToken для пагинации. Если объектов больше - нужен повторный запрос с continuation-token. При миллионе объектов в «директории» - 1000 HTTP-запросов только для листинга.
DELETE Object и DeleteObjects - удаление одного или батча объектов (до 1000 за раз):
POST /?delete HTTP/1.1
Host: bucket.s3.eu-central-1.amazonaws.com
Content-Type: application/xml
<Delete>
<Object><Key>warehouse/tmp/_temporary/part-001.parquet</Key></Object>
<Object><Key>warehouse/tmp/_temporary/part-002.parquet</Key></Object>
</Delete>
S3A использует batch delete при очистке _temporary/ директорий после commit - это существенно эффективнее, чем удалять по одному объекту.
AWS Signature V4: как подписываются запросы¶
Каждый HTTP-запрос к S3 должен быть подписан, чтобы S3 мог верифицировать идентичность отправителя. AWS Signature V4 - это алгоритм формирования подписи, который включает в себя:
- Canonical Request - нормализованное представление запроса (метод, URI, параметры, заголовки, хэш тела).
- String to Sign - строка из алгоритма подписи, даты, scope (регион + сервис) и хэша Canonical Request.
- Signing Key - ключ, полученный серией HMAC-SHA256 операций над Secret Key + дата + регион + сервис +
aws4_request. - Signature - HMAC-SHA256(Signing Key, String to Sign), преобразованный в hex.
Authorization: AWS4-HMAC-SHA256
Credential=AKIAIOSFODNN7EXAMPLE/20240116/eu-central-1/s3/aws4_request,
SignedHeaders=content-length;host;x-amz-content-sha256;x-amz-date,
Signature=fe5f80f77d5fa3beca038a248ff027d0445342fe2855ddc963176630326f1024
Почему это важно для Spark-инженера:
fs.s3a.signing-algorithmпереключает алгоритм между Signature V4 (современный) и V2 (устаревший, требуется старыми версиями Ceph). По умолчанию S3A использует V4.- Временные credentials (STS AssumeRole, IRSA) включают
x-amz-security-tokenзаголовок с session token - S3A обрабатывает это автоматически. - Signature V4 включает хэш тела запроса в подпись - это защищает от man-in-the-middle атак при передаче данных. Поэтому даже при HTTP (без TLS) данные защищены от подделки (но не от перехвата).
Внутренняя архитектура AWS S3¶
AWS не публикует детали внутреннего устройства S3, но на основе публикаций re:Invent и исследовательских работ известно следующее:
Распределение по Availability Zones. S3 хранит каждый объект в минимум трёх AZ одного региона. Это означает, что S3 Standard предоставляет 99.999999999% (11 девяток) durability - вероятность потери данных за год статистически исчезающе мала. Ни один HDFS-кластер не обеспечивает такую надёжность на уровне hardware failures.
Storage Nodes и разбиение объектов. Большие объекты (загруженные через Multipart Upload) физически хранятся разбитыми на части на разных storage nodes. Именно поэтому CompleteMultipartUpload - это не просто «сохранить файл», а «создать запись в метадате, которая указывает на разные storage locations».
Metadata Service. S3 хранит метаданные объектов (ключ, размер, ETag, ACL, версии) отдельно от данных в распределённой базе метаданных. Именно это делает HEAD запросы и ListObjectsV2 операциями против metadata service, а не storage nodes - поэтому LIST дешевле по bandwidth, но всё равно имеет latency.
Rate Limits и prefix-based scaling. S3 автоматически горизонтально масштабирует пропускную способность на основе префикса ключа:
- До 5 500 GET/HEAD запросов в секунду на один префикс
- До 3 500 PUT/COPY/POST/DELETE запросов в секунду на один префикс
Под «префиксом» в этом контексте понимается первые несколько символов ключа (примерно первые 6–8 символов). S3 автоматически шардирует горячие префиксы - если один префикс получает слишком много запросов, S3 разделяет его на несколько шардов.
Практические следствия для Spark:
При записи 1000 партиций в s3a://bucket/warehouse/table/year=2024/ все объекты имеют одинаковый длинный префикс. S3 видит это как нагрузку на один префикс и может throttle запросы (HTTP 503 Slow Down). Историческое решение - рандомизировать начало ключа (abc123/warehouse/...), но современный S3 сам справляется с распределением нагрузки по горячим префиксам в течение 30–60 минут после начала нагрузки.
Consistency Model: от eventual к strong¶
До декабря 2020 года S3 имел eventually consistent семантику для операций LIST и перезаписи объектов. Это создавало класс проблем, специфичных именно для Spark на S3:
Сценарий 1: Missing output files. Job A записывает 100 объектов, затем Job B читает директорию. До 2020 года LIST мог не вернуть все 100 объектов сразу - часть была «ещё не видна». Job B думала, что данных меньше.
Сценарий 2: Stale reads. Job перезаписывала объект (PUT). Следующий GET мог вернуть старую версию из-за кэширования в промежуточных узлах.
С декабря 2020 года AWS S3 гарантирует strong read-after-write consistency для всех операций: после успешного PUT, объект немедленно виден во всех GET и LIST запросах. Это устранило целый класс проблем. Magic Committer стал надёжно работать именно после этого изменения.
MinIO и Ceph. MinIO также гарантирует strong consistency. Ceph RGW - зависит от версии и конфигурации: старые версии (до Quincy) могли иметь eventual consistency на некоторых операциях LIST. Для production Ceph рекомендуется явно проверять гарантии consistency вашей версии.
URL форматы: virtual-hosted vs path-style¶
S3 поддерживает два формата URL для доступа к объектам:
Virtual-hosted style (дефолт для AWS S3):
https://my-bucket.s3.eu-central-1.amazonaws.com/warehouse/events/file.parquet
my-bucket.s3..... AWS S3 создаёт wildcard DNS-запись *.s3.amazonaws.com автоматически.
Path-style (требуется для MinIO, Ceph):
https://minio.internal:9000/my-bucket/warehouse/events/file.parquet
fs.s3a.path.style.access=true обязателен для on-premise хранилищ.
Примечание для AWS: с 2020 года AWS объявила deprecation path-style URLs и планирует его убрать - но реально это не произошло, и path-style по-прежнему поддерживается. Для AWS S3 лучше использовать virtual-hosted style (дефолт).
Lifecycle Rules и управление данными¶
S3 поддерживает автоматические политики жизненного цикла объектов - это важно для Data Lake:
import boto3
s3 = boto3.client("s3")
# Настройка Lifecycle Policy
lifecycle_config = {
"Rules": [
{
"ID": "abort-multipart-uploads",
"Status": "Enabled",
"Filter": {"Prefix": ""},
"AbortIncompleteMultipartUpload": {"DaysAfterInitiation": 7},
},
{
"ID": "expire-tmp-files",
"Status": "Enabled",
"Filter": {"Prefix": "_temporary/"},
"Expiration": {"Days": 1},
},
]
}
s3.put_bucket_lifecycle_configuration(
Bucket="my-data-lake",
LifecycleConfiguration=lifecycle_config
)
Первое правило - автоматически удаляет незавершённые Multipart Uploads через 7 дней. Это защита от «денежных утечек»: каждый незавершённый upload хранится и тарифицируется как обычные данные, но не виден в стандартном листинге.
Второе правило - удаляет объекты с префиксом _temporary/ через 1 день. Если Spark-джоба упала после partial commit и оставила файлы в _temporary/, они автоматически очистятся.
Как S3 масштабирует пропускную способность¶
Это один из наиболее важных разделов для понимания того, почему параллельные Spark-джобы иногда получают HTTP 503, почему порядок символов в ключе влияет на throughput и почему MinIO ведёт себя принципиально иначе, чем AWS S3.
Prefix Partitioning: внутренний механизм масштабирования S3¶
Когда AWS говорит о «лимитах 5 500 GET/s и 3 500 PUT/s на префикс», под словом «префикс» подразумевается не то, что большинство инженеров думают. Это не директория вида warehouse/events/ - это внутренний ключ шардирования S3, основанный на хэше первых символов имени объекта.
S3 хранит метаданные объектов в распределённой B-tree структуре, отсортированной лексикографически по ключу. Каждая нода этого дерева обслуживает конкретный диапазон ключей. Когда нагрузка на один диапазон вырастает выше порогового значения, S3 автоматически разбивает (splits) этот диапазон на два, каждый из которых обслуживается отдельным набором серверов.
Главный вывод: лимиты S3 не фиксированные - они применяются к одному partition в данный момент времени. По мере роста нагрузки S3 создаёт больше partitions, и суммарная пропускная способность бакета растёт пропорционально. При достаточно равномерном распределении нагрузки по ключам бакет может обслуживать сотни тысяч запросов в секунду.
Warm-up Period: S3 не масштабируется мгновенно¶
Ключевая проблема для Spark-джобов: split происходит не мгновенно. Когда новый бакет или новый префикс впервые получает высокую нагрузку, S3 обнаруживает перегрузку partition, принимает решение о split, реализует его - всё это занимает от 30 минут до нескольких часов.
В этот период «разогрева» (warm-up) запросы сверх лимита получают HTTP 503 Slow Down. S3A Connector делает retry с exponential backoff, но если нагрузка очень высокая и бакет холодный - часть запросов может падать.
Практическое следствие: первый запуск Spark-джобы, которая пишет данные в новый бакет с 500 Executor'ами, с высокой вероятностью получит волну 503 в первые 10–30 минут. После нескольких запусков (или после нескольких часов активной работы) бакет «разогреется» и throttling прекратится.
Стратегия для первого запуска на холодный бакет:
# 1. Начать с меньшего числа Executor'ов
spark = SparkSession.builder \
.config("spark.executor.instances", "20") \ # не 500
.config("spark.hadoop.fs.s3a.retry.limit", "15") \
.config("spark.hadoop.fs.s3a.retry.interval", "2000ms") \
.getOrCreate()
# 2. Записать небольшой объём данных - "прогреть" бакет
df_warmup = spark.range(0, 100_000)
df_warmup.write.parquet("s3a://new-bucket/warmup/")
# 3. После нескольких минут - основная нагрузка
# (бакет уже начал автоматически увеличивать партиции)
Дизайн ключей для максимального throughput¶
Поскольку S3 шардирует по хэшу первых символов ключа, структура имён объектов напрямую влияет на то, как нагрузка распределяется между partition'ами.
Плохая практика: дата в начале ключа
# Все объекты за 2024-01 начинаются с "2024/01/" → один S3 partition
# Вся нагрузка на один partition → быстрый throttling
s3a://bucket/2024/01/01/events/part-001.parquet
s3a://bucket/2024/01/01/events/part-002.parquet
s3a://bucket/2024/01/02/events/part-003.parquet
Ключи 2024/01/01/..., 2024/01/02/..., 2024/01/03/... все начинаются с 2024/01/ - лексикографически близкий диапазон. S3 держит их в одном-двух partition'ах.
Хорошая практика: рандомный или хэшированный префикс
# Ключи начинаются с хэша → нагрузка равномерно распределена
# UUID-prefix (Delta Lake / Iceberg делают это автоматически):
s3a://bucket/warehouse/a3f8bc2e-part-001.parquet
s3a://bucket/warehouse/d1e9fa7b-part-002.parquet
s3a://bucket/warehouse/0c45ab8f-part-003.parquet
# Или явное хэшированное партиционирование:
# <hash_prefix>/<дата>/<файл>
s3a://bucket/3f/2024/01/01/events/part-001.parquet # 3f = hash("part-001")[:2]
s3a://bucket/a7/2024/01/01/events/part-002.parquet # a7 = hash("part-002")[:2]
Хорошая новость для современных Lakehouse-форматов: Delta Lake и Apache Iceberg автоматически генерируют UUID-имена для файлов данных (part-00001-a3f8bc2e-d9c1-4f7e-b123-abc456def789-c000.snappy.parquet). Такие UUID равномерно распределены по всему алфавиту - S3 автоматически шардирует их по разным partition'ам. Именно поэтому Delta Lake/Iceberg практически не страдают от S3 throttling даже при высоком параллелизме.
Для Parquet без Lakehouse-формата - ключи вида part-00001.parquet, part-00002.parquet начинаются с одинакового part- → они попадают в один partition. При записи 1000 таких файлов одновременно из 200 Executor'ов - гарантированный throttling в первые минуты.
Расчёт реальной пропускной способности для Spark¶
Понимание масштабирования S3 позволяет рассчитать, сколько данных ваш кластер может реально записать/прочитать в секунду и при каком количестве Executor'ов начнётся throttling.
Теоретические лимиты AWS S3 (разогретый бакет):
Пропускная способность S3:
- Запись: ~10 Gbps на prefix (практически неограничено при UUID-ключах)
- Чтение: ~60 Gbps aggregate через parallel GET requests
При 200 Executor × 10 ядер × 1 Gbps сеть каждый = 200 Gbps суммарно
→ Сеть кластера быстрее S3 при плохих ключах; при хороших ключах S3 успевает
Практическая формула расчёта риска throttling:
# Расчёт нагрузки на S3
num_executors = 100
connections_per_executor = 200 # fs.s3a.connection.maximum
concurrent_requests = num_executors * connections_per_executor # = 20 000 req/s
# Лимит одного partition: 3500 PUT/s
# При холодном бакете с sequential ключами:
# 20 000 запросов → 1-2 partition'а → throttling гарантирован
# При UUID-ключах и разогретом бакете:
# 20 000 запросов → ~6 partition'ов → каждый получает ~3300 req/s → на грани
# Безопасная конфигурация: connection.maximum × executor_count ≤ 3500 × num_partitions
# Приблизительно: connection.maximum ≤ 100 на executor при 100 executor'ах
# (исходим из того, что S3 создаст 6-8 partition'ов на разогретом бакете)
print(f"Рекомендуемый connection.maximum: {3500 * 6 // num_executors}")
# → 210 на executor - что и является типичным оптимальным значением
Для MinIO: нет внутреннего auto-scaling. Пропускная способность ограничена железом: суммарный IOPS дисков, сетевая карта (10/25/100 Gbps). При превышении - не HTTP 503, а просто увеличение latency запросов и timeout.
Мониторинг throttling в Spark¶
Понять, испытывает ли ваша джоба throttling от S3, можно несколькими способами:
1. Spark UI → Stages → Tasks: колонка Task Duration
Если большинство тасок работают несколько секунд, но отдельные таски занимают 30–60 секунд при том же объёме данных - это признак retry из-за 503. S3A делает exponential backoff: первый retry через 500ms, второй через 1000ms, третий через 2000ms... При 7 retry суммарное ожидание может достигать 60+ секунд.
2. Spark-логи: поиск 503
# Фильтрация логов Spark Driver на 503 ошибки
grep "503\|SlowDown\|ServiceUnavailable\|Retrying" /path/to/spark-driver.log | head -20
# Пример строки лога при throttling:
# WARN S3ARetryPolicy: Retrying request to PUT s3a://bucket/warehouse/part-001.parquet
# after com.amazonaws.services.s3.model.AmazonS3Exception:
# Slow Down (Service: Amazon S3; Status Code: 503; Error Code: SlowDown)
3. Prometheus + Grafana: метрики S3A
При правильно настроенном hadoop-metrics2 S3A публикует счётчики retry:
# prometheus scrape config для Spark на Kubernetes
- job_name: spark_s3a
metrics_path: /metrics
static_configs:
- targets: ['spark-driver:4040']
Ключевая метрика: s3a_store_io_throttled_duration_ms - суммарное время, проведённое в ожидании retry. Рост этой метрики = сигнал о throttling.
MinIO: горизонтальное масштабирование через Erasure Coding¶
MinIO в distributed-режиме масштабируется принципиально иначе, чем AWS S3: нет автоматического partition split, зато есть линейное масштабирование через добавление серверов.
Erasure Coding в MinIO - это способ хранения данных с избыточностью без полной репликации. Вместо хранения 3 копий объекта (как HDFS с replication=3), MinIO разбивает объект на N+K частей, где:
- N - количество data shards (минимально нужных для восстановления)
- K - количество parity shards (избыточных для fault tolerance)
При схеме EC:8+4 (8 data + 4 parity, 12 нод суммарно): объект 1 GB разбивается на 12 частей по ~85 MB, размещённых на 12 нодах. Для чтения нужны любые 8 из 12 частей. Disk overhead = 1.5× (vs 3× для HDFS replication=3).
Следствие для Spark: при чтении одного объекта из MinIO EC:8+4 данные реально читаются с 8 нод одновременно - это дополнительный параллелизм на уровне хранилища. Пропускная способность одного GET-запроса ограничена не диском одной ноды, а суммарной пропускной способностью 8 нод.
Сравнение: как масштабируется S3 vs MinIO vs Ceph¶
| Аспект | AWS S3 | MinIO | Ceph RGW |
|---|---|---|---|
| Механизм масштабирования | Автоматический partition split | Линейное добавление нод | Добавление OSD (Ceph storage daemons) |
| Warm-up | 30 мин - несколько часов | Нет (linear scaling) | Нет |
| Лимит одного prefix | 3 500 PUT/s, 5 500 GET/s (до split) | Ограничен железом | Ограничен железом |
| Суммарный лимит | Практически неограничен при равномерных ключах | Линейно зависит от числа нод | Линейно зависит от OSD |
| Throttling под нагрузкой | HTTP 503 + exponential backoff | Рост latency + timeout | Рост latency + timeout |
| Избыточность данных | 3 AZ, 11 девяток durability | EC N+K или replication | EC или replication |
| Ключи для макс. throughput | UUID-образные, без общего prefix | Не критично | Не критично |
Эволюция Hadoop S3-коннекторов¶
Hadoop прошёл через три поколения коннекторов к S3:
s3:// - первое поколение (Hadoop 1.x). Использовал собственный MapReduce-native протокол, хранил метаданные блоков в специальных объектах. Несовместим с S3 REST API в современном понимании. Полностью устарел и удалён.
s3n:// (S3 Native) - второе поколение. Работал с реальным S3 REST API, но использовал старый HTTP-клиент JetS3t. Не поддерживал объекты крупнее 5 GB, имел проблемы с производительностью при многопоточном доступе и не поддерживал IAM-роли. Устарел, но встречается в legacy-коде.
s3a:// (S3 Advanced) - современное поколение, разработанное командой Apache Hadoop на базе официального AWS Java SDK. Поддерживает объекты до 5 TB (через Multipart Upload), нативные IAM-роли, конфигурируемые пулы потоков, Magic Committer. Используйте только s3a:// во всех новых пайплайнах.
Архитектура S3A Connector: Hadoop FileSystem Abstraction¶
Spark не работает напрямую с дисками, S3 или HDFS. Он использует Hadoop FileSystem API - унифицированный интерфейс, за которым скрывается любое хранилище. Когда вы пишете spark.read.parquet("s3a://bucket/path"), Spark:
- Смотрит на URI scheme (
s3a://) и ищет зарегистрированную реализациюFileSystemдля этой схемы - Вызывает методы этой реализации:
listStatus(),open(),create(),rename(),delete() - Получает данные в стандартизированном формате - InputStream или OutputStream
Благодаря этой абстракции Spark-код одинаково работает с HDFS, S3, Azure ADLS, Google GCS и MinIO - нужна только правильная реализация FileSystem. S3A Connector - это и есть реализация для всех S3-compatible хранилищ.
Важно понимать структуру S3A Connector: это не просто «прокси» к S3 API, а полноценная реализация файловой системы с собственным HTTP connection pool, движком параллельной загрузки (Multipart Upload Engine) и системой коммитов (Committer). Каждый из этих компонентов имеет свои параметры настройки.
Подключение S3A: зависимости и JAR¶
S3A Connector поставляется в составе hadoop-aws библиотеки. В кластерных дистрибутивах (CDH, HDP, EMR) она обычно уже включена в classpath. В standalone Spark или контейнерных средах её нужно добавить явно.
Добавление через spark-submit¶
# Версии hadoop-aws и aws-java-sdk-bundle должны строго совпадать
# с версией Hadoop, с которой скомпилирован ваш Spark
spark-submit \
--packages org.apache.hadoop:hadoop-aws:3.3.4,\
com.amazonaws:aws-java-sdk-bundle:1.12.367 \
--conf spark.hadoop.fs.s3a.endpoint=http://minio:9000 \
your_script.py
Критически важно: версия hadoop-aws должна совпадать с версией Hadoop в вашем Spark дистрибутиве. Проверить: spark.version и sc._jvm.org.apache.hadoop.util.VersionInfo.getVersion(). Несовпадение версий - самая частая причина NoSuchMethodError при работе с S3A.
Через SparkSession builder¶
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("S3A Demo") \
.config(
"spark.jars.packages",
"org.apache.hadoop:hadoop-aws:3.3.4,"
"com.amazonaws:aws-java-sdk-bundle:1.12.367"
) \
.getOrCreate()
Через кастомный Docker образ (Production рекомендация)¶
В production-среде правильнее включить нужные JAR в Docker-образ Spark:
FROM apache/spark:3.4.1
# Добавляем S3A JAR в папку с плагинами Spark
ADD https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/3.3.4/hadoop-aws-3.3.4.jar \
/opt/spark/jars/
ADD https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.12.367/aws-java-sdk-bundle-1.12.367.jar \
/opt/spark/jars/
Это надёжнее, чем --packages при каждом запуске: не требует интернета в runtime, версии зафиксированы в образе, время запуска джобы не тратится на скачивание JAR.
Аутентификация и безопасность¶
Провайдеры учётных данных (Credential Providers)¶
S3A поддерживает несколько способов аутентификации, расставленных в цепочке (Credential Provider Chain): если первый провайдер не возвращает credentials, пробуется следующий.
1. Access Key + Secret Key (SimpleAWSCredentialsProvider)
Самый простой способ - прямое указание ключей в конфигурации. Антипаттерн для production: ключи видны в spark-submit аргументах, логах, окружении.
spark = SparkSession.builder \
.config("spark.hadoop.fs.s3a.access.key", "AKIAIOSFODNN7EXAMPLE") \
.config("spark.hadoop.fs.s3a.secret.key", "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY") \
.getOrCreate()
2. Переменные окружения (EnvironmentVariableCredentialsProvider)
Лучше, чем явные ключи в коде. Переменные AWS_ACCESS_KEY_ID и AWS_SECRET_ACCESS_KEY автоматически подхватываются AWS SDK.
export AWS_ACCESS_KEY_ID=AKIAIOSFODNN7EXAMPLE
export AWS_SECRET_ACCESS_KEY=wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY
spark = SparkSession.builder \
.config(
"spark.hadoop.fs.s3a.aws.credentials.provider",
"com.amazonaws.auth.EnvironmentVariableCredentialsProvider"
) \
.getOrCreate()
3. IAM Roles / Instance Profile (Production стандарт для AWS)
На EC2-инстансах и в Kubernetes (IRSA - IAM Roles for Service Accounts) credentials получаются автоматически из метадата-сервиса AWS. Никаких ключей в коде, ротация происходит автоматически.
spark = SparkSession.builder \
.config(
"spark.hadoop.fs.s3a.aws.credentials.provider",
"com.amazonaws.auth.InstanceProfileCredentialsProvider"
) \
.getOrCreate()
Для Kubernetes с IRSA:
spark = SparkSession.builder \
.config(
"spark.hadoop.fs.s3a.aws.credentials.provider",
"com.amazonaws.auth.WebIdentityTokenCredentialsProvider"
) \
.getOrCreate()
4. Цепочка провайдеров (Production рекомендация)
В реальных системах лучше задать цепочку, где S3A последовательно пробует каждый провайдер:
spark = SparkSession.builder \
.config(
"spark.hadoop.fs.s3a.aws.credentials.provider",
",".join([
"com.amazonaws.auth.InstanceProfileCredentialsProvider",
"com.amazonaws.auth.EnvironmentVariableCredentialsProvider",
"com.amazonaws.auth.profile.ProfileCredentialsProvider",
])
) \
.getOrCreate()
Шифрование¶
Шифрование при передаче (in-transit): по умолчанию fs.s3a.connection.ssl.enabled=true - данные передаются по HTTPS. Для внутренних контуров с HTTP (dev/test MinIO без TLS) отключают: false.
Шифрование на стороне сервера (SSE - Server-Side Encryption): управляется параметром fs.s3a.server-side-encryption-algorithm:
# SSE-S3: AWS управляет ключами
.config("spark.hadoop.fs.s3a.server-side-encryption-algorithm", "AES256")
# SSE-KMS: ключи из AWS KMS
.config("spark.hadoop.fs.s3a.server-side-encryption-algorithm", "SSE-KMS")
.config("spark.hadoop.fs.s3a.server-side-encryption.key", "arn:aws:kms:...")
Проблема Rename и S3A Committers¶
Как Spark пишет данные: дефолтный FileOutputCommitter¶
Когда Spark выполняет df.write.parquet("s3a://bucket/table/"), он не пишет файлы сразу в финальную директорию. Он использует commit protocol:
- Каждый Executor пишет свою часть данных во временную папку:
s3a://bucket/table/_temporary/0/_temporary/attempt_xxx/part-00001.parquet - После успешного завершения всех Executor'ов Driver выполняет commit: переименовывает все файлы из
_temporary/в финальную директориюs3a://bucket/table/ - Временная папка удаляется
Зачем такая сложность? Это гарантирует атомарность на уровне job: если один из Executor'ов упал - его частичные данные остаются в _temporary/ и не «видны» читателям. Только полностью успешный job делает commit и делает данные видимыми.
На HDFS это работает отлично: rename - это O(1) операция в памяти NameNode. 1000 файлов - 1000 inode-обновлений, занимает миллисекунды.
На S3 это катастрофа: rename = copy + delete. 1000 файлов × 128 MB × 2 HTTP-запроса = ещё 256 GB сетевого трафика плюс 2000 HTTP-запросов только на commit phase. При 10 000 файлов commit занимает минуты.
FileOutputCommitter v1 vs v2¶
Hadoop предлагал две версии дефолтного коммитера, управляемые параметром mapreduce.fileoutputcommitter.algorithm.version:
Version 1 (дефолт): все временные файлы сначала собираются в одну папку _temporary/, затем Driver переименовывает их все за один проход. Безопаснее при сбоях, но commit phase - O(N) по числу файлов.
Version 2: каждый Executor переименовывает свои файлы сразу после завершения (не ждёт commit Driver'а). Быстрее, но если Driver упал после частичного commit - часть файлов уже видна читателям, часть ещё нет. Частичная видимость данных.
На S3 оба варианта плохи, потому что rename всё равно = copy+delete.
S3A Committers: решение проблемы¶
Команда Apache Hadoop разработала набор коммитеров, специально предназначенных для object storage. Они полностью обходят проблему rename.
Magic Committer¶
Magic Committer - самый производительный коммитер для AWS S3 (с 2020 года, когда появилась strong consistency). Вместо временных файлов + rename он использует S3 Multipart Upload:
- Executor инициирует Multipart Upload с уникальным ID
- Каждый чанк данных (обычно 64–128 MB) загружается как «part» - невидимый для читателей незавершённый Multipart
- После успешного завершения Executor записывает в «_magic» директорию манифест с ID multipart upload
- Driver читает все манифесты и вызывает
CompleteMultipartUploadдля каждого файла - атомарная операция, которая мгновенно «материализует» все загруженные части как полноценный объект
# Включение Magic Committer
spark = SparkSession.builder \
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.mapreduce.outputcommitter.factory.scheme.s3a",
"org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory") \
.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 с strong consistency. На MinIO с версии RELEASE.2020-01-03 - работает. На старых Ceph (< Quincy) - может не работать надёжно.
Directory Committer (Staging Committer)¶
Более консервативный вариант: Executor записывает данные во временную директорию на локальном диске воркера (не в S3!), а после успешного завершения таски загружает готовые файлы напрямую в финальную S3 директорию.
Плюсы: не требует strict consistency от S3, надёжно работает со старыми Ceph. Минус: нужно свободное место на локальных дисках Executor'ов.
# Directory Committer (более консервативный, работает со старыми Ceph)
spark = SparkSession.builder \
.config("spark.hadoop.fs.s3a.committer.name", "directory") \
.config("spark.hadoop.mapreduce.outputcommitter.factory.scheme.s3a",
"org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory") \
.config("spark.sql.sources.commitProtocolClass",
"org.apache.spark.internal.io.cloud.PathOutputCommitProtocol") \
.getOrCreate()
Коммитеры и Lakehouse-форматы¶
Если вы используете Delta Lake или Apache Iceberg - проблема коммитера практически снимается. Delta Lake пишет каждый файл с уникальным UUID-именем напрямую (без _temporary/), а транзакционная атомарность обеспечивается через ACID transaction log (_delta_log/). Файл либо записан в лог - и виден, либо нет. Rename не происходит вообще.
Для чистого Parquet без Lakehouse-формата выбор правильного коммитера критически важен.
| Формат записи | Нужен ли S3A Committer? |
|---|---|
| Parquet / ORC / CSV | Да, обязательно (Magic или Directory) |
| Delta Lake | Нет - собственный commit protocol |
| Apache Iceberg | Нет - object-store-native design |
| Apache Hudi | Нет - собственный механизм |
Multipart Upload: механика загрузки больших файлов¶
S3 ограничивает одиночный PUT-запрос размером 5 GB. Для больших файлов (а в Spark это типичный случай - 128 MB–2 GB Parquet) используется Multipart Upload:
CreateMultipartUpload- клиент инициирует multipart upload, получаетuploadIdUploadPart× N - загрузка каждого chunk'а (минимум 5 MB, максимум 5 GB) с номером частиCompleteMultipartUpload- клиент отправляет список всех частей, S3 «собирает» их в один объект атомарно- Альтернативно:
AbortMultipartUpload- если что-то пошло не так
Каждая часть загружается независимо и параллельно. S3A делает это в отдельном пуле потоков, что позволяет одному Executor'у загружать несколько chunks одновременно и максимально утилизировать сетевой канал.
Незавершённые Multipart Uploads - скрытая статья расходов в облаке. Если джоба упала в середине загрузки, незавершённые части остаются в S3 и списываются как хранение (по цене обычного хранения). AWS Console не показывает их в обычном листинге. Нужна политика жизненного цикла (Lifecycle Policy) или периодическая очистка:
import boto3
s3 = boto3.client("s3")
# Найти и удалить незавершённые multipart uploads
paginator = s3.get_paginator("list_multipart_uploads")
for page in paginator.paginate(Bucket="your-bucket"):
for upload in page.get("Uploads", []):
print(f"Abort upload: {upload['Key']}, started: {upload['Initiated']}")
s3.abort_multipart_upload(
Bucket="your-bucket",
Key=upload["Key"],
UploadId=upload["UploadId"]
)
Конфигурация производительности S3A¶
Connection Pool и Thread Pool¶
fs.s3a.connection.maximum (дефолт: 96) - максимальное количество одновременных HTTP-соединений с S3-сервером на один Executor. Каждый поток S3A при загрузке или скачивании занимает одно соединение. При параллельном чтении 200 партиций с одного Executor - дефолтных 96 соединений хватит, но если Executor читает и пишет одновременно (сложные pipeline) - может не хватить.
.config("spark.hadoop.fs.s3a.connection.maximum", "200")
fs.s3a.threads.max (дефолт: 96) - размер пула потоков для асинхронных S3-операций. Загрузка chunk'ов в Multipart Upload происходит в этом пуле. Увеличение до 200–500 на мощных воркерах позволяет параллелить больше операций.
.config("spark.hadoop.fs.s3a.threads.max", "200")
Важно: connection.maximum ≥ threads.max. Если потоков больше, чем соединений - часть потоков будет ждать свободного соединения, что убивает параллелизм.
Оптимизация Multipart Upload¶
fs.s3a.multipart.size (дефолт: 100M) - размер одного chunk при Multipart Upload. Оптимальные значения:
- 64 MB - для сетей с нестабильным качеством (меньше потерь при ретрансмиссии)
- 128 MB - стандарт для надёжных корпоративных сетей
- 256 MB+ - для очень быстрых (~10 Gbps) сетей с NVMe-дисками
.config("spark.hadoop.fs.s3a.multipart.size", "134217728") # 128 MB
fs.s3a.multipart.threshold (дефолт: 2147483647 - 2 GB) - порог размера объекта, выше которого включается Multipart Upload. Для Spark это практически означает «всегда single-part upload» при дефолтном значении! Снизьте до размера одного chunk:
# Включить multipart upload для файлов > 128 MB
.config("spark.hadoop.fs.s3a.multipart.threshold", "134217728")
fs.s3a.fast.upload (дефолт: true) - включение режима быстрой загрузки. В этом режиме данные буферизируются и загружаются параллельно в несколько потоков вместо последовательной записи.
fs.s3a.fast.upload.buffer - где буферизировать chunks перед отправкой. Критически важный выбор:
| Значение | Где буфер | Плюсы | Минусы |
|---|---|---|---|
disk |
Локальный диск Executor | Не занимает RAM | Доп. I/O на диск |
array |
JVM Heap | Быстро | Растёт куча, риск GC |
bytebuffer |
Off-heap (прямая память) | Быстро + нет GC | Требует executor.memoryOverhead |
# Для Executor с достаточным off-heap (spark.executor.memoryOverhead)
.config("spark.hadoop.fs.s3a.fast.upload.buffer", "bytebuffer")
# Для Executor с NVMe-дисками (надёжнее, нет риска OOM)
.config("spark.hadoop.fs.s3a.fast.upload.buffer", "disk")
Оптимизация чтения¶
fs.s3a.readahead.range (дефолт: 64K) - сколько байт S3A читает вперёд при открытии объекта. При чтении Parquet-файла Spark сначала читает footer (метаданные schema и row groups, обычно несколько MB в конце файла). С дефолтным readahead 64K это требует многих HTTP-запросов для получения footer'а. Увеличение до 1–4 MB резко ускоряет время открытия Parquet-файла:
.config("spark.hadoop.fs.s3a.readahead.range", "4194304") # 4 MB
fs.s3a.experimental.fadvise - режим чтения данных, аналог POSIX fadvise():
normal(дефолт) - универсальный режимsequential- оптимизация под линейное чтение (CSV, JSON). S3A агрессивно заранее читает следующие части файлаrandom- оптимизация под случайный доступ (Parquet, ORC). S3A не делает readahead, но кэширует части файла для повторного доступа
Для Parquet используйте random - Spark при column pruning читает только нужные column chunks, а не весь файл последовательно.
.config("spark.hadoop.fs.s3a.experimental.fadvise", "random")
Retry и Throttling¶
S3 и его аналоги имеют лимиты на запросы (Rate Limits). AWS S3 ограничивает: 5 500 GET/sec и 3 500 PUT/sec на префикс. При интенсивной параллельной записи из 100+ Executor'ов можно получить HTTP 503 Slow Down.
fs.s3a.retry.limit (дефолт: 7) - сколько раз повторять запрос при transient ошибках (503, сетевые таймауты).
fs.s3a.retry.interval (дефолт: 500ms) - базовый интервал между retry. S3A использует exponential backoff с jitter.
.config("spark.hadoop.fs.s3a.retry.limit", "10")
.config("spark.hadoop.fs.s3a.retry.interval", "1000ms")
fs.s3a.attempts.maximum (дефолт: 20) - максимальное число попыток для AWS SDK внутри одного HTTP-запроса.
Paging операций LIST¶
fs.s3a.paging.maximum (дефолт: 5000) - количество объектов, возвращаемых за один LIST-запрос (максимум S3 API - 1000 для ListObjects, 1000 для ListObjectsV2). Параметр управляет batch size при листинге тысяч партиций. Уменьшать не нужно; максимум определяется API S3.
Специфика on-premise: MinIO и Ceph¶
При работе с on-premise S3-совместимым хранилищем (MinIO, Ceph RGW) требуется дополнительная конфигурация, которая не нужна при работе с AWS S3.
Обязательные параметры для MinIO¶
fs.s3a.endpoint - URL MinIO-сервера вместо дефолтного s3.amazonaws.com. Указывает на балансировщик MinIO или конкретную ноду:
.config("spark.hadoop.fs.s3a.endpoint", "http://minio.internal:9000")
# Или для HTTPS с кастомным сертификатом:
.config("spark.hadoop.fs.s3a.endpoint", "https://minio.internal:9443")
fs.s3a.path.style.access - критически важный параметр. AWS S3 использует virtual-hosted style: https://bucket.s3.amazonaws.com/key. MinIO и Ceph используют path style: https://minio.host:9000/bucket/key. Без path.style.access=true Spark попытается сделать DNS-запрос для bucket.minio.internal - и получит ошибку.
.config("spark.hadoop.fs.s3a.path.style.access", "true")
fs.s3a.connection.ssl.enabled - отключить SSL для HTTP-эндпоинтов MinIO в dev/test окружении:
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false")
fs.s3a.impl - явное указание реализации FileSystem:
.config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")
Специфика Ceph RGW¶
Ceph RGW (RADOS Gateway) - S3-совместимый шлюз для Ceph. Некоторые версии Ceph имеют особенности:
fs.s3a.signing-algorithm - алгоритм подписи запросов. Старые Ceph (< Luminous) могут требовать v2 подпись вместо современной v4:
# Для старых Ceph, не поддерживающих AWS Signature V4
.config("spark.hadoop.fs.s3a.signing-algorithm", "S3SignerType")
Таймауты для нагруженного Ceph. Если Ceph находится под высокой фоновой нагрузкой (scrubbing, rebalancing) - запросы могут медленно отвечать:
.config("spark.hadoop.fs.s3a.connection.establish.timeout", "10000") # 10s
.config("spark.hadoop.fs.s3a.connection.timeout", "200000") # 200s
Production-ready конфигурация¶
Полный шаблон для MinIO (on-premise)¶
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("OpenSource-DataLake-ETL") \
\
.config("spark.hadoop.fs.s3a.endpoint", "http://minio.internal:9000") \
.config("spark.hadoop.fs.s3a.access.key", "your-access-key") \
.config("spark.hadoop.fs.s3a.secret.key", "your-secret-key") \
.config("spark.hadoop.fs.s3a.path.style.access", "true") \
.config("spark.hadoop.fs.s3a.impl",
"org.apache.hadoop.fs.s3a.S3AFileSystem") \
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false") \
\
.config("spark.hadoop.fs.s3a.connection.maximum", "200") \
.config("spark.hadoop.fs.s3a.threads.max", "200") \
.config("spark.hadoop.fs.s3a.multipart.size", "134217728") \
.config("spark.hadoop.fs.s3a.multipart.threshold", "134217728") \
.config("spark.hadoop.fs.s3a.fast.upload", "true") \
.config("spark.hadoop.fs.s3a.fast.upload.buffer", "disk") \
.config("spark.hadoop.fs.s3a.readahead.range", "4194304") \
.config("spark.hadoop.fs.s3a.experimental.fadvise","random") \
\
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.mapreduce.outputcommitter.factory.scheme.s3a",
"org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory") \
.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.retry.limit", "10") \
.config("spark.hadoop.fs.s3a.retry.interval", "1000ms") \
\
.getOrCreate()
Полный шаблон для AWS S3 (с IAM Role)¶
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("AWS-DataLake-ETL") \
\
.config(
"spark.hadoop.fs.s3a.aws.credentials.provider",
"com.amazonaws.auth.InstanceProfileCredentialsProvider,"
"com.amazonaws.auth.EnvironmentVariableCredentialsProvider"
) \
\
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "true") \
.config("spark.hadoop.fs.s3a.server-side-encryption-algorithm", "AES256") \
\
.config("spark.hadoop.fs.s3a.connection.maximum", "200") \
.config("spark.hadoop.fs.s3a.threads.max", "200") \
.config("spark.hadoop.fs.s3a.multipart.size", "134217728") \
.config("spark.hadoop.fs.s3a.multipart.threshold", "134217728") \
.config("spark.hadoop.fs.s3a.fast.upload", "true") \
.config("spark.hadoop.fs.s3a.fast.upload.buffer", "bytebuffer") \
.config("spark.hadoop.fs.s3a.readahead.range", "4194304") \
.config("spark.hadoop.fs.s3a.experimental.fadvise","random") \
\
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.mapreduce.outputcommitter.factory.scheme.s3a",
"org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory") \
.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()
S3 и Data Locality: почему locality всегда ANY¶
При работе с HDFS Spark scheduler размещает таски на тех Executor'ах, которые физически находятся на том же сервере, что и нужные блоки данных (уровни PROCESS_LOCAL, NODE_LOCAL, RACK_LOCAL).
При работе с object storage (S3, MinIO) data locality почти всегда ANY. Причины:
- Object Storage - это отдельный кластер серверов (в облаке - вообще другая инфраструктура)
- Данные не привязаны к конкретным узлам (у S3 нет понятия «блок на конкретном сервере» с точки зрения клиента)
- Spark Executor'ы не знают, на каких физических дисках хранится конкретный object
Это означает, что каждый read из S3 проходит через сеть. Вся данные передаются по сети между object storage и Spark Executor'ами. При плохой сетевой связности или перегруженном storage это становится узким местом.
Компенсация:
- Высокопроизводительная сеть (10 Gbps, 25 Gbps, 100 Gbps в современных дата-центрах)
- Параллельное чтение из множества потоков (настройки connection pool)
- Predicate Pushdown и Column Pruning - читать только нужные данные
- Правильный размер файлов: крупные Parquet-файлы = меньше HTTP-соединений = меньше overhead
Мониторинг и диагностика¶
Spark UI метрики для S3¶
В Spark UI → Stage → Tasks смотрите колонки:
- Shuffle Read Size/Records - сетевой трафик при shuffle (не связан с S3, но влияет на общую нагрузку сети)
- Input Size/Records - объём прочитанных данных из S3. При колонном pruning должен быть значительно меньше физического размера файлов
- Task Duration - если duration задач высокий при малом Input Size → bottleneck в S3 latency (много мелких HTTP-запросов)
Hadoop метрики S3A через JMX¶
S3A Connector публикует метрики через Hadoop metrics framework. Включить их в Spark:
spark = SparkSession.builder \
.config("spark.hadoop.hadoop.metrics2.impl",
"org.apache.hadoop.metrics2.impl.MetricsSystemImpl") \
.getOrCreate()
# Метрики доступны через SparkUI → Environment → Hadoop Properties
# Или через Prometheus если настроен hadoop-metrics2-prometheus
Ключевые метрики S3A:
S3AInstrumentation.STORE_IO_REQUEST_DURATION- время HTTP-запросов к S3S3AInstrumentation.OP_LIST_OBJECTS- количество LIST операций (высокое = проблема мелких файлов)S3AInstrumentation.MULTIPART_UPLOAD_STARTED- сколько Multipart Upload инициированоS3AInstrumentation.MULTIPART_UPLOAD_ABORTED- сколько прервано (должно быть 0)
Диагностика через логи¶
# Включить DEBUG-логирование S3A (очень подробно, только для диагностики)
spark-submit \
--conf "spark.driver.extraJavaOptions=-Dlog4j.logger.org.apache.hadoop.fs.s3a=DEBUG" \
your_script.py
# В логах будет видно:
# DEBUG S3AFileSystem: Opening 's3a://bucket/file.parquet'
# DEBUG S3AFileSystem: GET https://minio:9000/bucket/file.parquet bytes=[0-65535]
# DEBUG S3AFileSystem: Completed GET in 45ms, 65536 bytes
Типичные ошибки и их причины¶
| Ошибка | Причина | Решение |
|---|---|---|
Connection refused: minio:9000 |
Неверный endpoint | Проверить fs.s3a.endpoint |
DNS resolution failed: bucket.minio.internal |
Не включён path-style access | fs.s3a.path.style.access=true |
403 Forbidden |
Неверные credentials или нет прав | Проверить access key/secret или IAM policy |
HTTP 503 Slow Down |
Rate limiting S3 | Увеличить fs.s3a.retry.limit, уменьшить параллелизм |
MultipartUploadError: EntityTooSmall |
Chunk < 5 MB | fs.s3a.multipart.size ≥ 5 MB |
SignatureDoesNotMatch |
Неверный алгоритм подписи | Для старого Ceph: fs.s3a.signing-algorithm=S3SignerType |
ClassNotFoundException: S3AFileSystem |
Нет hadoop-aws JAR | Добавить через --packages или в Docker-образ |
Практика: полноценный пайплайн Spark + MinIO¶
Развёртывание MinIO через Docker Compose¶
# docker-compose.yml
version: "3.8"
services:
minio:
image: minio/minio:RELEASE.2024-01-01T00-00-00Z
command: server /data --console-address ":9001"
environment:
MINIO_ROOT_USER: minioadmin
MINIO_ROOT_PASSWORD: minioadmin123
ports:
- "9000:9000" # S3 API
- "9001:9001" # Web Console
volumes:
- minio_data:/data
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:9000/minio/health/live"]
interval: 30s
timeout: 20s
retries: 3
volumes:
minio_data:
docker compose up -d
# Создать bucket через MinIO CLI
docker run --rm --network host \
minio/mc:latest alias set local http://localhost:9000 minioadmin minioadmin123
docker run --rm --network host \
minio/mc:latest mb local/data-lake
docker run --rm --network host \
minio/mc:latest ls local/
PySpark ETL с MinIO: от сырых данных до Parquet¶
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
# ─────────────────────────────────────────────────────────────
# 1. Инициализация SparkSession с полной конфигурацией MinIO
# ─────────────────────────────────────────────────────────────
spark = SparkSession.builder \
.appName("Spark-MinIO-ETL") \
.master("local[*]") \
.config("spark.jars.packages",
"org.apache.hadoop:hadoop-aws:3.3.4,"
"com.amazonaws:aws-java-sdk-bundle:1.12.367") \
\
.config("spark.hadoop.fs.s3a.endpoint", "http://localhost:9000") \
.config("spark.hadoop.fs.s3a.access.key", "minioadmin") \
.config("spark.hadoop.fs.s3a.secret.key", "minioadmin123") \
.config("spark.hadoop.fs.s3a.path.style.access", "true") \
.config("spark.hadoop.fs.s3a.impl",
"org.apache.hadoop.fs.s3a.S3AFileSystem") \
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false") \
\
.config("spark.hadoop.fs.s3a.multipart.size", "67108864") \ # 64 MB
.config("spark.hadoop.fs.s3a.multipart.threshold", "67108864") \
.config("spark.hadoop.fs.s3a.fast.upload", "true") \
.config("spark.hadoop.fs.s3a.fast.upload.buffer", "disk") \
.config("spark.hadoop.fs.s3a.readahead.range", "1048576") \ # 1 MB
.config("spark.hadoop.fs.s3a.experimental.fadvise","random") \
\
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.mapreduce.outputcommitter.factory.scheme.s3a",
"org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory") \
.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()
# ─────────────────────────────────────────────────────────────
# 2. Генерируем тестовые данные (Bronze layer)
# ─────────────────────────────────────────────────────────────
df_raw = spark.range(0, 10_000_000).select(
F.col("id"),
F.date_sub(F.current_date(), (F.col("id") % 365).cast("int")).alias("event_date"),
(F.col("id") % 50).alias("category_id"),
F.rand().alias("amount"),
F.when(F.col("id") % 100 == 0, None).otherwise(F.col("id") * 1.5).alias("revenue")
)
# Записываем сырые данные как Parquet в Bronze-слой MinIO
df_raw.coalesce(8).write \
.mode("overwrite") \
.parquet("s3a://data-lake/bronze/events/")
print("Bronze layer записан")
# ─────────────────────────────────────────────────────────────
# 3. Silver layer: очистка, типизация, партиционирование
# ─────────────────────────────────────────────────────────────
df_bronze = spark.read.parquet("s3a://data-lake/bronze/events/")
print(f"Читаем Bronze: {df_bronze.count()} строк, {df_bronze.rdd.getNumPartitions()} партиций")
df_silver = df_bronze \
.filter(F.col("revenue").isNotNull()) \
.withColumn("event_year", F.year("event_date")) \
.withColumn("event_month", F.month("event_date"))
# Партиционирование по году и месяцу для partition pruning
df_silver.write \
.partitionBy("event_year", "event_month") \
.mode("overwrite") \
.parquet("s3a://data-lake/silver/events/")
print("Silver layer записан")
# ─────────────────────────────────────────────────────────────
# 4. Проверяем partition pruning: читаем только нужную партицию
# ─────────────────────────────────────────────────────────────
df_filtered = spark.read.parquet("s3a://data-lake/silver/events/") \
.filter(F.col("event_year") == 2024) \
.filter(F.col("event_month") == 1)
# В explain должны быть PartitionFilters, а не просто PushedFilters
df_filtered.explain(mode="formatted")
# Запрос читает только s3a://data-lake/silver/events/event_year=2024/event_month=1/
# Остальные партиции не трогает - LIST запросы идут только на нужную партицию
Сравнение: дефолтный коммитер vs Magic Committer¶
import time
def benchmark_write(spark, df, path, use_magic_committer=True, label=""):
if use_magic_committer:
spark.conf.set("spark.hadoop.fs.s3a.committer.name", "magic")
spark.conf.set("spark.sql.sources.commitProtocolClass",
"org.apache.spark.internal.io.cloud.PathOutputCommitProtocol")
else:
# Дефолтный FileOutputCommitter
spark.conf.set("spark.sql.sources.commitProtocolClass",
"org.apache.spark.sql.execution.datasources.SQLHadoopMapReduceCommitProtocol")
start = time.time()
df.write.mode("overwrite").parquet(path)
elapsed = time.time() - start
print(f"{label}: {elapsed:.1f}s")
return elapsed
df_test = spark.range(0, 5_000_000)
t_default = benchmark_write(
spark, df_test, "s3a://data-lake/bench/default/",
use_magic_committer=False, label="FileOutputCommitter (default)"
)
t_magic = benchmark_write(
spark, df_test, "s3a://data-lake/bench/magic/",
use_magic_committer=True, label="Magic Committer"
)
print(f"\nМагический коммитер быстрее в {t_default / t_magic:.1f}x")
# Типичный результат на MinIO: Magic Committer быстрее в 2-5x при 200+ файлах
Антипаттерны при работе с S3¶
1. Дефолтный FileOutputCommitter на S3 для Parquet. Rename-фаза превращает каждый commit в O(N) copy+delete операций. При 10 000 файлах - минуты только на финализацию.
2. Хранение миллионов мелких файлов. Каждый файл = отдельный HTTP-запрос при листинге. S3 LIST возвращает 1000 объектов за запрос. 1 млн файлов = 1000 LIST запросов просто чтобы прочитать список. Planification времени в Spark растёт пропорционально.
3. fs.s3a.connection.maximum слишком высокий на MinIO. На AWS S3 лимиты практически бесконечны. На ваш корпоративный MinIO-сервер сотня Executor'ов по 200 соединений = 20 000 одновременных соединений. Это может быть DoS-атакой на ваш же сервер. Для MinIO балансируйте: executor_count × connection.maximum ≤ capacity MinIO.
4. Хранение credentials в открытом виде в SparkSession.builder или spark-defaults.conf. Эти конфиги видны в Spark UI (Environment tab), в логах кластера, в git-истории. Используйте переменные окружения или Kubernetes Secrets.
5. Отсутствие политики очистки незавершённых Multipart Uploads. После падения джобы незавершённые части остаются и оплачиваются. Настройте Lifecycle Policy на удаление multipart uploads старше 7 дней или запустите периодическую очистку через boto3.
6. fadvise=sequential для Parquet. При column pruning Spark читает не весь файл последовательно, а прыгает по row groups. Sequential prefetching скачивает данные впустую. Для Parquet - random.
7. Игнорирование версионной совместимости hadoop-aws и aws-java-sdk. Неправильная пара версий даёт NoSuchMethodError или ClassNotFoundException в runtime - самые сложные для диагностики ошибки. Всегда проверяйте матрицу совместимости.
Сравнительная таблица: AWS S3 vs MinIO vs Ceph¶
| Параметр | AWS S3 | MinIO | Ceph RGW |
|---|---|---|---|
fs.s3a.endpoint |
Не нужен (дефолт) | http://host:9000 |
http://host:7480 |
fs.s3a.path.style.access |
false (дефолт) |
true |
true |
fs.s3a.connection.ssl.enabled |
true (дефолт) |
false (dev) / true (prod) |
Зависит от конфигурации |
fs.s3a.signing-algorithm |
Не нужен | Не нужен | S3SignerType (старые версии) |
| Magic Committer | Рекомендован | Работает (> 2020) | Осторожно (проверить версию) |
| Strong Consistency | Да (с 2020) | Да | Зависит от версии Ceph |
| Credentials | IAM Role / ключи | Ключи / LDAP | Ключи |
| Rate Limits | 5500 GET/s, 3500 PUT/s | Ограничено железом | Ограничено железом |
Итоговый чек-лист¶
- Используйте только
s3a://-s3://иs3n://устарели - Включите Magic Committer для Parquet без Lakehouse-формата - иначе commit = rename = copy+delete
- Delta Lake / Iceberg не нуждаются в специальных коммитерах - у них собственный ACID commit
fs.s3a.path.style.access=trueобязателен для MinIO, Ceph и всех on-prem S3-совместимых хранилищfs.s3a.multipart.threshold = fs.s3a.multipart.size- включает Multipart Upload для реальных размеров файловfs.s3a.fast.upload.buffer=disk- надёжный выбор для большинства конфигураций без риска OOMfs.s3a.experimental.fadvise=random- обязательно для Parquet с column pruningfs.s3a.readahead.range=1-4 MB- ускоряет чтение Parquet footer'аconnection.maximum = threads.maxи оба ≥ количество одновременных Executor-задач- Очищайте незавершённые Multipart Uploads - они незаметно накапливаются и стоят денег
- Файлы 128 MB–1 GB - оптимальный размер для S3; мелкие файлы убивают performance листинга
- credentials никогда не в коде и git - переменные окружения, Kubernetes Secrets или IAM Role