Сравнение S3-совместимых хранилищ: MinIO, Ceph, SeaweedFS, Garage, Apache Ozone
Детальный технический разбор пяти open-source S3-совместимых хранилищ для Data Lake и Lakehouse: архитектурные требования Spark, внутренняя анатомия каждого движка, сравнение производительности метаданных, интеграция с Delta Lake и Iceberg, тюнинг SparkConf и траблшутинг.
1. Архитектурные требования Spark к On-Premise S3-хранилищам¶
Прежде чем сравнивать конкретные продукты, необходимо понять, что именно Spark требует от объектного хранилища. Требования аналитических Big Data нагрузок принципиально отличаются от требований типичного веб-приложения, которое хранит картинки пользователей и выгружает их по одной.
Чем Big Data нагрузка отличается от веб-нагрузки¶
Веб-приложение работает с S3 через паттерн «один пользователь - один объект»: GET одного файла, PUT одного файла. Запросы независимы, небольшие по размеру, случайны по ключам. Для такой нагрузки любой S3-совместимый движок подойдёт.
Spark-пайплайн работает совершенно иначе:
- Тысячи параллельных операций одновременно. При чтении партиционированной таблицы 200 Executor'ов одновременно отправляют LIST-запросы к хранилищу, затем параллельно читают сотни Parquet-файлов. Хранилище получает взрывной трафик, который длится минуты.
- Тяжёлые последовательные чтения. Parquet Row Group размером 128–512 MB читается целиком одним запросом. Производительность определяется пропускной способностью сети и диска, а не задержкой.
- Миллисекундные SLA на операции с метаданными. LIST, HEAD объекта - это критический путь планировщика Spark. Если LIST 10 000 файлов занимает 30 секунд, всё это время Executor'ы простаивают.
- Атомарная финализация записи. Spark пишет данные во временные файлы, затем «публикует» их через операцию rename или commit. Некорректная реализация этой операции в хранилище может привести к видимости частичных данных.
Схема показывает, что Spark нагружает хранилище на двух принципиально разных уровнях: метаданный (LIST, HEAD - latency-чувствительный) и данных (GET, PUT больших блоков - throughput-чувствительный). Хранилище, которое хорошо справляется с одним, может провалиться на другом.
LIST Latency: узкое горлышко планирования¶
LIST - самая критичная операция для Spark. Прежде чем читать данные, Spark должен перечислить все файлы в партиции. Это происходит на Driver в однопоточном режиме, то есть не параллелизуется.
Пример: таблица с 10 000 файлами Parquet, разбитыми по датам за 3 года. Spark Driver вызывает FileSystem.listStatus(), который транслируется в серию LIST-запросов к хранилищу. Если каждый LIST-запрос возвращает 1000 объектов (стандартное ограничение S3 API), нужно 10 запросов подряд. При latency 20 мс на запрос - 200 мс. При latency 200 мс (нагруженный Ceph RGW) - 2 секунды. Это ещё терпимо.
Но если таблица содержит 1 000 000 файлов (частый результат стриминга без компакции) - разница между latency 20 мс и 200 мс на один LIST-запрос выливается в 20 секунд vs 200 секунд только на перечисление файлов. И всё это время Executor'ы простаивают.
Строгая консистентность: read-after-write¶
Когда Spark-таск записывает файл в хранилище, другой таск (или сам Driver) должен сразу же его увидеть при LIST и GET. Это называется strong read-after-write consistency.
AWS S3 достиг этого только в 2020 году. До этого S3 был eventual consistent: только что записанный объект мог не появиться в LIST мгновенно. Spark S3A Connector содержал целый слой workaround'ов для борьбы с этой проблемой.
Для on-premise хранилищ это требование проще выполнить, потому что нет географически распределённой репликации. Но некоторые системы (особенно с eventual consistency моделью) всё равно могут нарушать это требование при высокой нагрузке.
Иллюзия директорий и проблема атомарного Rename¶
S3 - это плоское пространство ключей, а не иерархическая файловая система. «Директории» в S3 - это просто общий префикс в имени ключа. s3a://bucket/data/2024/01/15/file.parquet - это один объект с ключом data/2024/01/15/file.parquet, а не файл в иерархии директорий.
Проблема: POSIX-семантика (на которой строился Hadoop FileSystem API) предполагает атомарную операцию rename. Spark использует rename для финализации записи: сначала записывает данные в _temporary/, затем переименовывает в финальный путь. На HDFS rename - атомарная операция, занимающая миллисекунды.
На S3 нативного rename не существует. Вместо него выполняется COPY + DELETE. Для директории с тысячами файлов это:
- PUT каждого файла по новому ключу (одновременно или последовательно)
- DELETE каждого файла по старому ключу
Это долго, дорого и не атомарно: в промежутке между шагами 1 и 2 читатель может увидеть оба варианта файла. Именно поэтому были созданы специальные Spark Committer'ы (Magic, Staging, Directory), которые обходят проблему rename по-разному. Поведение конкретного хранилища напрямую влияет на то, какой Committer работает корректно.
2. MinIO и Garage: Shared-Nothing архитектуры для средних масштабов¶
MinIO и Garage объединяет общая философия: минимальная операционная сложность, отсутствие выделенных серверов метаданных и ориентация на self-hosted развёртывание. Но их целевые сценарии кардинально отличаются.
Анатомия MinIO: Go, Erasure Coding и встроенный consensus¶
MinIO написан на Go и распространяется как единственный статически слинкованный бинарник без внешних зависимостей. Это одно из главных архитектурных решений: отсутствие ZooKeeper, etcd, Consul и других систем координации, которые в Ceph или HDFS являются отдельными компонентами.
Ключевой архитектурный факт: в MinIO нет выделенной базы метаданных. Метаданные каждого объекта хранятся в файле xl.meta рядом с данными на каждом диске. xl.meta - это небольшой JSON-файл, содержащий:
- Erasure Coding параметры (какой шард на каком диске)
- Checksums каждого шарда
- Пользовательские метаданные объекта
- Версии объекта (если включено версионирование)
Это решение имеет важное следствие: LIST в MinIO - это обход файловой системы. Для перечисления объектов с префиксом MinIO должен прочитать xl.meta с достаточного числа дисков, чтобы выполнить кворум. При миллионах объектов это становится узким местом: тысячи мелких операций чтения метаданных с дисков.
Производительность LIST в MinIO (типичные значения):
- 10 000 объектов в одном префиксе: ~50-200 мс
- 100 000 объектов: ~500 мс - 2 сек
- 1 000 000 объектов: ~5-20 сек (зависит от дисков и числа нод)
На NVMe дисках производительность существенно лучше: случайное чтение миллионов мелких xl.meta файлов - именно то, в чём NVMe выигрывает у HDD на несколько порядков.
Magic Committer и MinIO. MinIO реализует специальный механизм «magic files» - возможность записывать объекты в специальный путь, делая их невидимыми до явного commit'а. Это позволяет Spark Magic Committer'у работать без staging-директории и операции COPY+DELETE. Результат: запись Spark в MinIO с Magic Committer на ~30-40% быстрее, чем с Directory Committer.
Слабые стороны MinIO для больших масштабов¶
MinIO имеет архитектурное ограничение: добавление нод только целыми Erasure Set'ами. Нельзя добавить одну ноду к кластеру из 4 нод - нужно добавить ещё 4 ноды (или 8, 16 - кратно исходному Erasure Set). Это неудобно при постепенном росте кластера.
Более серьёзная проблема при масштабировании: метаданные на дисках не индексируются. При миллиардах объектов LIST-операции деградируют, потому что нет централизованного индекса. В сравнении с Ceph (который хранит метаданные в специальных пулах с B-tree индексами) или Apache Ozone (RocksDB-based OM) - это ограничение.
Анатомия Garage: Rust, Kademlia DHT и geo-distributed consensus¶
Garage - принципиально другой подход. Он написан на Rust и спроектирован для развёртывания на гетерогенном железе в нескольких географических зонах с нестабильными сетевыми связями между ними.
Kademlia DHT вместо RAFT. В отличие от MinIO (RAFT требует кворума живых нод) и Ceph (PAXOS-подобный консенсус в Monitor кластере), Garage использует вариант Kademlia Distributed Hash Table для координации. Это означает, что:
- Нет «лидера», который мог бы стать узким местом
- Кластер работает при потере произвольного подмножества нод, если хотя бы одна нода с репликой доступна
- Операции eventual consistent - нет гарантии, что только что записанный объект мгновенно виден на всех нодах
Почему Garage не подходит для аналитического Spark-пайплайна:
Для Data Lake пишут паттерн: Spark записывает файлы, затем Driver делает commit (видимость данных). При этом предполагается, что после commit другие Spark-джобы видят новые данные. В eventual consistency системе это не гарантировано. Интервал сходимости (convergence) Garage - от секунд до минут в зависимости от сетевой задержки между зонами.
Для Spark ETL это неприемлемо: пайплайн B, запускаемый сразу после завершения пайплайна A, может не увидеть данные, записанные пайплайном A.
Где Garage реально полезен: хранение файлов веб-приложений, резервные копии, static assets в геораспределённой CDN-подобной архитектуре. Для аналитики с высокими требованиями к консистентности - не лучший выбор.
Garage под нагрузкой: хаос-тестирование в боевых условиях¶
Независимое сравнение Garage и SeaweedFS (версии v2.1 и v4.07 соответственно) в условиях инъекции сетевых сбоев выявило важные поведенческие различия, которые не видны на синтетических бенчмарках в идеальных условиях.
Базовый throughput при параллельной загрузке 2 GB данных:
- SeaweedFS: 31 MB/s - в 2.1 раза выше базового показателя Garage
- Garage: 14.5 MB/s
Разрыв объясняется архитектурно: SeaweedFS записывает данные напрямую в Volume-файлы с батчингом, тогда как Garage выполняет синхронную репликацию на все ноды перед подтверждением записи клиенту.
Поведение при сетевых сбоях:
При имитации задержки 100 мс на двух нодах одновременно Garage прерывал тест из-за connection timeout. Аналогичный сценарий с одной нодой давал деградацию до 13.5 MB/s - кластер продолжал работать, но медленнее. SeaweedFS при задержке 100 мс на одной ноде деградировал до 10.8 MB/s, но не терял соединения.
При потере 2% пакетов на двух нодах одновременно Garage показывал 10 MB/s и 228 секунд на 2 GB. SeaweedFS при аналогичном сценарии восстанавливался за 3 секунды и продолжал работу с минимальной деградацией.
Накладные расходы на метаданные:
Garage хранит около 2.5 KB служебных метаданных на каждый объект - для тестового набора из 12 000 файлов по 256 KB это 30 MB только метаданных. SeaweedFS по заявленным характеристикам тратит около 40 байт на файл (in-memory Volume index), хотя практические данные зависят от настроек Filer backend.
Восстановление после отказа зоны - наиболее критичный сценарий для Data Lake. Garage при потере ноды требует ручного запуска garage repair blocks после восстановления. Без этого кластер считает себя консистентным, хотя часть блоков может быть недореплицирована. SeaweedFS обнаруживал отсутствующие реплики автоматически и запускал восстановление в фоне без вмешательства оператора.
Эти результаты подтверждают вывод из предыдущего раздела: Garage - интересная система для геораспределённых сценариев с толерантностью к eventual consistency, но для аналитических Data Lake пайплайнов, где важны предсказуемость и автоматическое восстановление, он требует дополнительной операционной осторожности.
3. Ceph RGW и Apache Ozone: Enterprise-гиганты для экзабайтных масштабов¶
Ceph и Apache Ozone - системы для совершенно другого масштаба и уровня операционной зрелости. Оба продукта рассчитаны на большие инженерные команды, многолетний жизненный цикл и сотни-тысячи нод.
Анатомия Ceph: RADOS и многослойная архитектура¶
Ceph - это не просто объектное хранилище. Это распределённая storage платформа, поверх которой реализованы три независимых интерфейса: RGW (S3/Swift), RBD (блочные устройства) и CephFS (файловая система). Все три используют единый нижний слой - RADOS (Reliable Autonomic Distributed Object Store).
CRUSH Algorithm - ключевое отличие Ceph от других систем. CRUSH (Controlled Replication Under Scalable Hashing) - это детерминированная функция, которая по ключу объекта и текущей Cluster Map вычисляет, на каких OSD должны лежать реплики (или шарды Erasure Code). Клиент вычисляет это локально, без обращения к центральному серверу метаданных. Это позволяет Ceph масштабироваться на тысячи OSD без деградации производительности операций данных.
BlueStore - хранилище данных нового поколения. В современных версиях Ceph данные на OSD хранятся не в файловой системе ОС (как в старом FileStore), а в специальном движке BlueStore, который работает напрямую с блочным устройством. BlueStore даёт:
- Полный контроль над I/O без overhead'а VFS Linux
- Встроенные checksums для обнаружения bit rot
- Нативную поддержку Erasure Coding без двойной записи
Metadata Pools в Ceph. Для S3-объектов через RGW Ceph хранит метаданные в специальных пулах (pools). Пул - это логическое разделение объектов внутри кластера с собственными правилами репликации. Метаданные RGW хранятся в:
.rgw.root- конфигурация самого RGW.rgw.control- синхронизация RGW инстансов.rgw.meta- пользовательские данные, buckets.rgw.buckets.index- индексы содержимого бакетов
.rgw.buckets.index - самый важный пул для Spark. Именно в нём хранится индекс объектов каждого бакета. LIST-операции обращаются именно туда. При большом числе объектов в одном бакете этот пул становится узким местом. Рекомендуется размещать его на быстрых NVMe-дисках и выделять достаточное число PG (Placement Groups).
Операционная реальность Ceph¶
Ceph требует специализированного инженера - Ceph Administrator. Это не преувеличение. Ключевые задачи, которые нельзя автоматизировать без глубокого понимания системы:
- CRUSH Map tuning. Правила распределения реплик по failure domain (rack, host, disk). Неправильный CRUSH Map означает, что три реплики одного объекта окажутся на одном rack'е - и потеря rack'а потеряет данные.
- PG (Placement Groups) расчёт. Слишком мало PG - горячие PG становятся bottleneck. Слишком много PG - overhead на metadata.
- OSD recovery tuning. После потери диска Ceph восстанавливает данные, интенсивно нагружая оставшиеся OSD. Неправильные настройки recovery throttle могут деградировать производительность на дни.
Эта сложность - не недостаток Ceph, а цена за его гибкость и масштабируемость.
Анатомия Apache Ozone: Cloud-Native замена HDFS¶
Apache Ozone появился из наблюдения: HDFS при хранении миллиардов файлов упирается в ограничение NameNode. NameNode хранит метаданные всех файлов в Java heap. Практический предел - несколько сотен миллионов файлов на крупный кластер.
Ozone решает это через разделение на два сервиса с разными функциями.
OM (Ozone Manager) хранит метаданные на диске через RocksDB. В отличие от HDFS NameNode, которая требует загрузить всё в RAM при старте (что занимает часы на больших кластерах), OM читает метаданные с диска по запросу. Теоретический предел - сотни миллиардов объектов.
SCM (Storage Container Manager) управляет Container'ами - физическими единицами хранения размером 5 GB. Container содержит много Chunk'ов (блоков данных из разных объектов). Важно: SCM управляет Container'ами, а не отдельными объектами. Это существенно снижает metadata overhead по сравнению с HDFS (где NameNode управляла каждым блоком).
Два клиентских протокола - стратегически важное решение. Ozone поддерживает:
- OzoneFileSystem (ofs://) - нативный Hadoop FileSystem API. Spark подключается без S3A Connector, напрямую через Ozone JARы. Это быстрее, потому что нет HTTP overhead'а и нет трансляции через S3 Gateway.
- S3 Gateway - HTTP шлюз для совместимости с существующим кодом, написанным для AWS S3. Тот же код, что работает с MinIO, работает с Ozone через S3 Gateway без изменений.
Наличие нативного OzoneFileSystem - главное конкурентное преимущество Ozone перед MinIO и Ceph для Hadoop-экосистемы.
4. SeaweedFS: Специализированное решение проблемы Small Files¶
SeaweedFS занимает уникальную нишу: это не универсальное объектное хранилище, а система, специально спроектированная для сценария, с которым другие системы справляются плохо - хранения и быстрой раздачи огромного числа мелких объектов.
Вдохновение: Facebook Haystack¶
Facebook в середине 2000-х годов столкнулся с проблемой: социальная сеть хранила миллиарды фотографий пользователей. При хранении в обычной файловой системе (NFS) каждый GET фотографии требовал нескольких дисковых операций:
- Прочитать inode директории
- Прочитать inode файла
- Прочитать данные файла
Для маленьких файлов операции 1 и 2 - pure overhead. Facebook создал Haystack: систему, которая упаковывает тысячи мелких объектов в один большой файл (Superblock) и поддерживает in-memory индекс offset'ов. GET фотографии требует одной операции чтения: offset из индекса + seeked read в Superblock.
SeaweedFS реализует ту же концепцию в open source.
Архитектура SeaweedFS¶
Volume = большой файл-контейнер. Каждый Volume Server хранит объекты в Volume-файлах фиксированного размера (по умолчанию 30 GB). Внутри Volume файла объекты упакованы последовательно. Рядом с каждым .dat файлом Volume'а - .idx файл с оффсетами объектов. Поиск объекта: находим Volume через Master (in-memory таблица), читаем оффсет из .idx (последовательный доступ), читаем данные из .dat (seeked read по известному оффсету).
Для мелких файлов (KB) это принципиально быстрее чем отдельный inode в ext4/XFS: нет random I/O по дереву директорий файловой системы.
Разделение namespace-семантики и хранения. Filer Layer отвечает за понятие «директорий» и S3-совместимый API. Он хранит маппинг path → FileID во внешнем key-value хранилище. Это критически важное архитектурное решение: metadata throughput масштабируется независимо от throughput хранения данных. Если LIST работает медленно - заменяем Filer backend с LevelDB на Cassandra или TiKV.
Почему SeaweedFS полезен для Bronze-слоя Spark Streaming¶
Structured Streaming и Kafka Consumer записывают данные в S3 в микробатчах каждые 30–120 секунд. Каждый микробатч создаёт несколько небольших Parquet-файлов (1–50 MB). За 24 часа накапливается 700–2880 файлов. За месяц - ~70 000 файлов. За год - ~800 000 файлов на один топик.
При таком паттерне накопления мелких файлов:
- MinIO: LIST 800 000 файлов работает медленно из-за обхода xl.meta файлов на дисках
- SeaweedFS: Filer хранит маппинг в RocksDB с индексом по префиксу - LIST работает быстро вне зависимости от числа объектов
После накопления за день данные компактируются в крупные Parquet-файлы и перекладываются в MinIO для Silver/Gold слоёв. SeaweedFS занимает роль быстрого приёмника сырых данных, а не долгосрочного хранилища.
Compaction в SeaweedFS: автоматическое уплотнение¶
После удаления объектов внутри Volume-файла остаются «мёртвые» байты. Volume размером 30 GB может содержать 50% живых данных и 50% удалённых. SeaweedFS выполняет Volume Compaction: перезаписывает живые объекты в новый Volume-файл, удаляет старый. Это аналог VACUUM в Delta Lake, но на уровне физического хранилища.
Compaction можно запустить вручную или настроить автоматически по расписанию. Во время compaction Volume помечается как read-only, новые записи идут в другие Volume'ы.
Операционные ограничения SeaweedFS, которые не очевидны с первого взгляда¶
Несмотря на впечатляющую производительность на мелких файлах, у SeaweedFS есть несколько практических ограничений, которые становятся заметны при промышленной эксплуатации.
Внешний бэкенд метаданных Filer обязателен при масштабировании. По умолчанию Filer использует встроенный LevelDB для хранения namespace маппинга (path → FileID). LevelDB - однопоточная embedded база данных, которая хорошо подходит для разработки, но не для production с высокой конкурентной нагрузкой. При интенсивном стриминге через Structured Streaming, когда сотни потоков одновременно пишут мелкие файлы, LevelDB становится bottleneck. Рекомендуемая конфигурация для production: RocksDB или внешняя база данных (PostgreSQL, MariaDB, Cassandra, TiKV). Это добавляет операционную сложность: необходимо поддерживать отдельный кластер метаданных.
Отсутствие нативного S3 Lifecycle Policy. SeaweedFS не реализует стандартный S3 Lifecycle API, позволяющий задавать правила автоматического удаления объектов по возрасту или тегам. Это означает, что задачи очистки устаревших данных (например, удаление raw JSON файлов из Bronze после компакции в Parquet) нужно реализовывать отдельно: либо через Airflow DAG с boto3, либо вручную через weed shell. Для Ceph и MinIO такие политики настраиваются декларативно через стандартный S3 API.
Незавершённые Multipart Upload не очищаются автоматически. В S3-совместимых системах незавершённые MPU (например, от упавшего Spark Job) занимают дисковое пространство и могут накапливаться. MinIO и Ceph поддерживают Lifecycle Rule AbortIncompleteMultipartUpload, которая автоматически удаляет такие «мусорные» загрузки. В SeaweedFS аналога нет - необходимо периодически запускать ручную команду s3.clean.uploads.
Kubernetes Operator vs Helm Chart. SeaweedFS имеет поддерживаемый Kubernetes Operator, который управляет rolling updates, масштабированием и recovery. Garage предоставляет только Helm Chart без Operator'а - операционные задачи приходится выполнять вручную.
Эти ограничения не делают SeaweedFS плохим выбором - они просто смещают его в более специализированную нишу: быстрый приёмник сырых данных с последующей компакцией и переносом в более зрелое хранилище. Пара SeaweedFS (Bronze) + MinIO (Silver/Gold) устраняет большинство проблем за счёт разделения ответственности.
5. Сравнительный батл: метаданные, Erasure Coding и TCO¶
Бенчмарк метаданных: что происходит при 200 параллельных Executor'ах¶
Рассмотрим конкретный сценарий: Spark читает партиционированную таблицу. Driver запускает обнаружение файлов (FileSystem.listStatus), затем 200 Executor'ов параллельно начинают читать файлы.
Важная оговорка: числа в диаграмме - ориентировочные. Реальная производительность зависит от числа нод, типа дисков (HDD vs SSD vs NVMe), сетевой конфигурации и количества объектов в одном префиксе. Эти числа - типичный порядок величин для кластеров из 4-8 нод на SSD.
Сравнительные цифры throughput: крупные объекты¶
Для аналитической нагрузки (крупные Parquet-файлы, Sequential GET/PUT) производительность на больших объектах принципиальна. Измерения на однотипном стенде (8-ядерный сервер, NVMe-диски, конфигурации от 4 до 8 нод) дают следующую картину:
| Хранилище | Конфигурация защиты | Чтение | Запись |
|---|---|---|---|
| MinIO | EC:4+4 (8 дисков) | ~2.8 GB/s | ~2.1 GB/s |
| SeaweedFS | Репликация 2x | ~2.3 GB/s | ~1.8 GB/s |
| Ceph RGW | EC:3+1 | ~1.9 GB/s | ~1.4 GB/s |
| Garage | Репликация 3x | ~1.6 GB/s | ~1.2 GB/s |
Обратите внимание: Ceph проигрывает MinIO по throughput не из-за слабого движка, а из-за дополнительного сетевого хопа через RGW-шлюз. RADOS напрямую (через librbd или librados) может давать производительность выше, чем через HTTP REST API. Для Spark это не применимо - S3A Connector всегда идёт через HTTP - но полезно понимать, откуда берётся этот overhead.
Задержка на малых объектах (файлы до 1 MB, latency P99):
| Хранилище | Средняя задержка PUT/GET (мс) |
|---|---|
| SeaweedFS | ~2.1 |
| MinIO | ~3.8 |
| Garage | ~4.2 |
| Ceph RGW | ~6.3 |
SeaweedFS здесь лидирует именно благодаря Volume-архитектуре: PUT мелкого объекта - это append в уже открытый Volume-файл с обновлением in-memory индекса. Ceph проигрывает из-за накладных расходов на RGW → librados → OSD цепочку.
Почему Garage медленнее на метаданных? Garage использует eventual consistency - при HEAD-запросе нода может не иметь самой свежей информации и должна дополнительно синхронизироваться. В single-datacenter развёртывании это несущественно, но накладывает дополнительный latency по сравнению с системами с strong consistency.
Эффективность Erasure Coding: сколько дискового пространства уходит «в воздух»¶
Erasure Coding - способ обеспечить отказоустойчивость с меньшим overhead'ом по сравнению с полной репликацией. Но разные системы поддерживают разные конфигурации.
| Хранилище | Схема защиты | Overhead дисков | Допустимые потери |
|---|---|---|---|
| MinIO | EC:4+4 (4 data + 4 parity) | 100% (2x) | любые 4 диска из 8 |
| MinIO | EC:6+2 (6 data + 2 parity) | 33% | любые 2 диска из 8 |
| Ceph | Erasure Code 4+2 | 50% | любые 2 OSD |
| Ceph | Репликация 3x | 200% | любые 2 реплики |
| SeaweedFS | Репликация 2x | 100% | 1 Volume Server |
| Apache Ozone | RS 3+2 | 67% | любые 2 DataNodes |
| Apache Ozone | Репликация 3x | 200% | любые 2 DataNodes |
| Garage | Репликация 3x (по зонам) | 200% | 2 зоны из 3 |
Что это означает практически. Если у вас 100 TB сырых дисков:
- MinIO с EC:6+2: ~75 TB полезного пространства (100 / 1.33)
- Ceph с EC 4+2: ~67 TB полезного пространства
- Репликация 3x (любой движок): ~33 TB полезного пространства
Для Data Lake с большим объёмом данных разница между 33 TB и 75 TB полезного пространства из одного и того же железа - это существенная экономия на капекс.
Потребление ресурсов самими хранилищами¶
Хранилище потребляет CPU и RAM даже в idle - для housekeeping, rebalancing, metadata caching. Это ресурсы, которые не достаются Spark.
| Хранилище | RAM на ноду | CPU idle (%) | Зависимости |
|---|---|---|---|
| MinIO | 256 MB – 2 GB | 1-3% | Нет (один бинарник) |
| Ceph OSD | 1-4 GB | 5-15% | Нет внешних, но несколько сервисов |
| Ceph MON | 2-8 GB | 2-5% | - |
| SeaweedFS | 512 MB – 4 GB | 2-5% | Опционально: etcd/RocksDB |
| Garage | 64-256 MB | <1% | Нет |
| Apache Ozone OM | 4-16 GB | 5-10% | ZooKeeper |
| Apache Ozone DataNode | 2-8 GB | 3-8% | - |
Garage потребляет наименьше ресурсов - это следствие минималистичного дизайна на Rust. Идеален для развёртывания на маломощных серверах или ARM-платах.
Apache Ozone OM потребляет много RAM из-за RocksDB кешей. Для кластера с миллиардами объектов OM рекомендуют развёртывать на серверах с 32-64 GB RAM.
Сводная инженерная матрица выбора¶
| Критерий | MinIO | Ceph RGW | SeaweedFS | Garage | Apache Ozone |
|---|---|---|---|---|---|
| Полнота S3 API | ★★★★★ | ★★★★☆ | ★★★☆☆ | ★★★☆☆ | ★★★☆☆ |
| LIST latency | ★★★★☆ | ★★★☆☆ | ★★★★★ | ★★★☆☆ | ★★★★☆ |
| Throughput больших объектов | ★★★★★ | ★★★★☆ | ★★★☆☆ | ★★★★☆ | ★★★★☆ |
| Small files performance | ★★★☆☆ | ★★☆☆☆ | ★★★★★ | ★★★☆☆ | ★★★☆☆ |
| Строгая консистентность | ✅ | ✅ | ✅ | ❌ | ✅ |
| Magic Committer support | ✅ | ❌ | ❌ | ❌ | ❌ |
| Сложность операций | Низкая | Очень высокая | Средняя | Очень низкая | Высокая |
| Масштабируемость | До ~100 PB | Практически неограничена | До ~10 PB | До ~1 PB | До ~10 PB |
| Порог входа (человеко-недели) | 0.5 | 4-8 | 1-2 | 0.5 | 2-4 |
| Hadoop-интеграция | ★★★☆☆ | ★★★☆☆ | ★★★☆☆ | ★★☆☆☆ | ★★★★★ |
| Потребление RAM (на ноду) | Низкое | Среднее | Среднее | Минимальное | Высокое |
Объективный тест совместимости S3 API¶
Звёздочки в матрице выше - качественная оценка. Существуют и количественные измерения: команда Vitastor project прогнала официальный S3-совместимый тест-сьют против основных открытых реализаций, чтобы выбрать фронтенд для своего хранилища.
Методология: запускается батарея тест-кейсов, покрывающих весь спектр S3 API - от базовых PUT/GET до Versioning, Object Lock, SSE-шифрования, Lifecycle, STS и CORS. Результат для каждой системы - число пройденных, упавших тестов и ошибок:
| Реализация | Язык | Пройдено | Упало | Ошибки | Итог |
|---|---|---|---|---|---|
| Ceph RadosGW | C++ | 576 | 69 | 34 | Лучшая совместимость |
| Zenko CloudServer | Node.js | 382 | 253 | 47 | Неожиданно высокий результат |
| MinIO | Go | 321 | 311 | 52 | Средняя совместимость |
| SeaweedFS | Go | 56 | 176 | 510 | Наихудший результат |
Несколько важных выводов из этих цифр.
Ceph - лидер S3-совместимости. 576 пройденных тестов - это существенно больше, чем у любой другой open source реализации. Ceph RGW разрабатывается с явной целью максимально точно воспроизвести поведение AWS S3, включая граничные случаи. Для enterprise-развёртываний, где существующий код написан с расчётом на строгое соответствие AWS S3, Ceph - самый безопасный выбор.
MinIO прошёл 321 тест - хуже, чем кажется. Это неожиданный результат, учитывая репутацию MinIO как «наиболее AWS-совместимого». Причина - MinIO намеренно не реализует ряд устаревших или редко используемых S3-операций: BucketLogging, Website hosting, PostObject (загрузка через HTML-форму), часть ACL-операций. Для Spark и стандартных Data Engineering задач это не проблема - все нужные операции (PUT, GET, LIST, Multipart, Versioning) MinIO реализует корректно. Но если ваш стек использует экзотические S3-возможности - проверьте совместимость заранее.
SeaweedFS: 56 из 500+ - красный флаг. Основная причина катастрофического результата - практически полное отсутствие Versioning, а также SSE-C, SSE-KMS, Object Lock, Lifecycle, CORS, STS, BucketPolicy. 510 ошибок (не просто упавших тестов, а именно ошибок - исключений и неожиданных HTTP статусов) указывают на несовместимость на уровне протокола в части API. Это подтверждает вывод: SeaweedFS - не универсальная S3-замена, а специализированный инструмент для конкретных сценариев. Тщательно проверяйте, нужны ли вашему пайплайну функции, которых нет.
Zenko CloudServer (Node.js) - неожиданно сильный результат. 382 пройденных теста - второе место, лучше MinIO. Это open source S3-шлюз от Scality, изначально разработанный для работы поверх различных бэкендов. Менее известен в Data Engineering сообществе, но заслуживает внимания для сценариев, где важна широкая S3-совместимость без операционной сложности Ceph.
6. Особенности интеграции с Lakehouse-форматами: Delta Lake и Apache Iceberg¶
Delta Lake и Apache Iceberg - это не просто форматы файлов. Это протоколы управления транзакциями поверх объектных хранилищ. Их корректная работа критически зависит от семантики операций хранилища.
Как Delta Lake использует объектное хранилище¶
Delta Lake хранит данные в Parquet-файлах и управляет транзакциями через директорию _delta_log/. Каждый commit создаёт новый JSON-файл в _delta_log/: 00000000000000000001.json, 00000000000000000002.json и т.д.
Ключевое требование к хранилищу: атомарность PUT-операции. Delta Lake протокол предполагает, что после успешного PUT файл 00002.json либо виден полностью, либо не виден вообще. Не должно быть состояния, когда файл виден частично (например, видна его «директория», но сам файл ещё пишется).
Все пять рассматриваемых хранилищ обеспечивают атомарность PUT - объект становится видимым только после полного завершения PUT-запроса. Это стандартное требование S3 API, и ни одна из систем его не нарушает.
Конкурентная запись: Optimistic Concurrency Control в Delta Lake¶
Delta Lake использует Optimistic Concurrency Control (OCC): несколько writer'ов пытаются коммитить одновременно, и один из них получает конфликт.
# Псевдокод Delta Lake Optimistic Concurrency
def commit_transaction(version_to_write):
# Шаг 1: записываем данные (Parquet файлы)
write_parquet_files(data)
# Шаг 2: пытаемся записать лог транзакций
log_file = f"_delta_log/{version_to_write:020d}.json"
# Шаг 3: проверяем, что версия ещё не занята
if object_exists(log_file):
# Кто-то другой уже закоммитил эту версию
# Читаем новый state, проверяем conflicts
current_version = get_latest_version()
if can_retry_at_version(current_version + 1):
return commit_transaction(current_version + 1) # retry
else:
raise ConcurrentModificationException()
# Шаг 4: атомарный PUT - ключевая операция
put_object(log_file, transaction_data)
Почему eventual consistency опасна для Delta Lake. Если хранилище имеет eventual consistency (например, Garage при межзональной репликации), то после успешного PUT файла 00002.json другой клиент может не увидеть его при HEAD-проверке. Это означает, что два writer'а одновременно посчитают, что версия 00002 свободна, и оба попытаются записать 00002.json. В результате один перезапишет коммит другого, и данные будут повреждены.
Именно поэтому для Delta Lake и Iceberg обязательна строгая read-after-write consistency. Из пяти рассматриваемых систем только Garage не гарантирует это в geo-distributed развёртывании.
Apache Iceberg: FileIO и уменьшение зависимости от LIST¶
Apache Iceberg спроектирован более «S3-friendly», чем Delta Lake. Одно из ключевых отличий - способ обнаружения файлов.
Delta Lake при открытии таблицы выполняет LIST директории _delta_log/, чтобы найти все commit-файлы. При большом числе транзакций (тысячи коммитов без checkpoint) это дорого.
Apache Iceberg использует metadata.json - манифестный файл, который содержит точный список всех Parquet-файлов таблицы. Spark читает один metadata.json, из него узнаёт пути ко всем файлам без LIST-операций.
Iceberg metadata structure:
s3a://warehouse/
├── my_table/
│ ├── metadata/
│ │ ├── 00000-abc.metadata.json ← точка входа (snapshot metadata)
│ │ ├── 00001-def.metadata.json ← после первой записи
│ │ ├── snap-123456789.avro ← snapshot manifest list
│ │ └── manifest-xyz.avro ← список Parquet файлов
│ └── data/
│ ├── 00000-0-part-0.parquet
│ └── 00000-1-part-0.parquet
Это означает, что LIST latency хранилища имеет меньшее значение для Iceberg, чем для Delta Lake. Даже если LIST медленный (как в Ceph при большом числе объектов) - Iceberg его почти не использует в hot path.
Зато атомарность PUT важна ещё больше. Iceberg использует собственный механизм оптимистичной конкуренции через обновление указателя metadata.json. Если PUT файла metadata.json не атомарен - таблица может оказаться в некорректном состоянии.
Iceberg FileIO: прямая работа с S3 API без Hadoop S3A¶
Традиционно Spark обращается к хранилищу через Hadoop FileSystem API (S3A Connector). Это слой абстракции, добавляющий overhead: Hadoop FileSystem не является нативным S3 клиентом и содержит POSIX-ориентированные паттерны (rename, директории), которые плохо подходят для объектных хранилищ.
Apache Iceberg предлагает альтернативу - FileIO: набор интерфейсов для работы с файлами, которые имплементируются под конкретное хранилище.
# Iceberg с S3FileIO (обходит Hadoop S3A)
spark = SparkSession.builder \
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
.config("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog") \
.config("spark.sql.catalog.local.type", "hadoop") \
.config("spark.sql.catalog.local.warehouse", "s3a://warehouse/") \
# Используем S3FileIO вместо HadoopFileIO
# S3FileIO работает через AWS SDK напрямую, минуя Hadoop FileSystem API
# Это даёт: меньше overhead, лучшая поддержка multipart, корректный ETag
.config("spark.sql.catalog.local.io-impl", "org.apache.iceberg.aws.s3.S3FileIO") \
# S3FileIO конфигурируется через iceberg.* параметры
.config("spark.sql.catalog.local.s3.endpoint", "http://minio:9000") \
.config("spark.sql.catalog.local.s3.path-style-access", "true") \
.config("spark.sql.catalog.local.s3.access-key-id", "minioadmin") \
.config("spark.sql.catalog.local.s3.secret-access-key", "minioadmin") \
.getOrCreate()
Преимущества S3FileIO для производительности:
- Нет overhead'а от POSIX-семантики Hadoop FileSystem
- Нет попыток эмулировать rename через COPY+DELETE
- Прямое использование Multipart Upload API для больших файлов
- Корректная обработка ETag для checksumming
Практический вывод: для хранилищ с неполной поддержкой Hadoop FileSystem API (Garage, SeaweedFS) - Iceberg с S3FileIO может работать лучше, чем через стандартный S3A Connector. S3FileIO требует только базового S3 API (PUT, GET, HEAD, LIST, DELETE), не используя специфические расширения.
7. Практика: тюнинг SparkConf и траблшутинг¶
Универсальный шаблон SparkSession для On-Premise S3¶
Следующий фрагмент - production-ready шаблон, который работает с MinIO, Ceph RGW, SeaweedFS и Apache Ozone (через S3 Gateway). Каждый параметр прокомментирован с объяснением «почему».
import os
from pyspark.sql import SparkSession
def create_spark_for_onprem_s3(
endpoint_url: str,
access_key: str,
secret_key: str,
app_name: str = "spark-app",
enable_magic_committer: bool = False,
max_connections: int = 200,
) -> SparkSession:
"""
Создаёт SparkSession для работы с On-Premise S3-совместимым хранилищем.
enable_magic_committer=True только для MinIO.
Для Ceph, SeaweedFS, Garage, Ozone - оставляйте False.
max_connections: число параллельных HTTP-соединений от одного Executor'а.
При 200 Executor'ах × 200 соединений = 40 000 одновременных соединений.
Убедитесь, что хранилище выдержит такую нагрузку (file descriptors, threads).
По умолчанию в Hadoop S3A = 200. Для нагруженных кластеров можно снизить до 50.
"""
builder = SparkSession.builder \
.appName(app_name) \
\
.config(
"spark.jars.packages",
"org.apache.hadoop:hadoop-aws:3.3.4,"
"com.amazonaws:aws-java-sdk-bundle:1.12.367"
) \
\
# ── Базовая конфигурация S3A Connector ───────────────────────────────
# endpoint: перенаправляем все s3a:// запросы в наш on-premise хост
.config("spark.hadoop.fs.s3a.endpoint", endpoint_url) \
\
# path.style.access: обязательно для on-premise эмуляторов.
# Virtual-Hosted style (по умолчанию) требует DNS wildcard вида
# *.storage.example.com - сложно настроить на on-premise.
.config("spark.hadoop.fs.s3a.path.style.access", "true") \
\
# credentials.provider: явно указываем провайдер.
# Без этого S3A пробует несколько провайдеров последовательно,
# включая InstanceProfileCredentialsProvider - он ждёт ответа от
# EC2 Instance Metadata Service 30 сек, которого нет on-premise.
.config(
"spark.hadoop.fs.s3a.aws.credentials.provider",
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
) \
.config("spark.hadoop.fs.s3a.access.key", access_key) \
.config("spark.hadoop.fs.s3a.secret.key", secret_key) \
\
# ssl.enabled: on-premise хранилища чаще всего без TLS (или с self-signed).
# false: не пытаться установить TLS соединение.
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false") \
\
# ── Производительность соединений ────────────────────────────────────
# connection.maximum: размер HTTP connection pool на каждый JVM процесс.
# Spark создаёт несколько JVM (один Driver + N Executor'ов).
# Итого: (1 + N) × max_connections соединений к хранилищу.
.config("spark.hadoop.fs.s3a.connection.maximum", str(max_connections)) \
\
# threads.max: размер thread pool для async операций S3A.
# Должен быть ≥ connection.maximum / 2.
.config("spark.hadoop.fs.s3a.threads.max", str(max_connections // 2)) \
\
# multipart.size: размер каждой части при Multipart Upload.
# 128 MB: оптимально для Parquet файлов 128-512 MB.
# Меньше → больше частей → больше API запросов.
# Больше → меньше частей, но больше памяти на буферизацию.
.config("spark.hadoop.fs.s3a.multipart.size", "134217728") \
\
# fast.upload: использовать буферизацию в памяти при записи.
# true (default): Spark буферизует данные в памяти перед PUT.
# Лучше для потоковой записи без прохода через временный файл на диске.
.config("spark.hadoop.fs.s3a.fast.upload", "true") \
\
# ── Retry логика ─────────────────────────────────────────────────────
# На нагруженных on-premise хранилищах иногда возникают transient ошибки.
# Автоматический retry с exponential backoff снижает число failed jobs.
.config("spark.hadoop.fs.s3a.retry.limit", "7") \
.config("spark.hadoop.fs.s3a.retry.interval", "500ms") \
.config("spark.hadoop.fs.s3a.attempts.maximum", "5") \
\
# ── Committer конфигурация ────────────────────────────────────────────
# Выбираем 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"
)
# Magic Committer - только для MinIO
if enable_magic_committer:
builder = builder \
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.fs.s3a.committer.magic.enabled", "true")
else:
# Directory Committer - универсальный, работает везде
# Staging создаёт staging директорию на стороне Spark (локально или S3)
# Directory: staging директория в целевом S3 prefix
builder = builder \
.config("spark.hadoop.fs.s3a.committer.name", "directory")
return builder.getOrCreate()
# ── Примеры использования ────────────────────────────────────────────────────
# MinIO (рекомендуется Magic Committer)
spark_minio = create_spark_for_onprem_s3(
endpoint_url="http://minio.example.com:9000",
access_key="minioadmin",
secret_key="minioadmin",
enable_magic_committer=True,
)
# Ceph RGW
spark_ceph = create_spark_for_onprem_s3(
endpoint_url="http://ceph-rgw.example.com:7480",
access_key="AKIAIOSFODNN7EXAMPLE",
secret_key="wJalrXUtnFEMI/K7MDENG",
enable_magic_committer=False, # Ceph не поддерживает Magic Committer
)
# SeaweedFS Filer
spark_seaweed = create_spark_for_onprem_s3(
endpoint_url="http://seaweedfs-filer.example.com:8333",
access_key="any_key",
secret_key="any_secret",
max_connections=50, # SeaweedFS чувствительнее к параллельным соединениям
enable_magic_committer=False,
)
# Apache Ozone S3 Gateway
spark_ozone = create_spark_for_onprem_s3(
endpoint_url="http://ozone-s3g.example.com:9878",
access_key="ozone_key",
secret_key="ozone_secret",
enable_magic_committer=False,
)
Диагностика типичных ошибок интеграции¶
При первом подключении Spark к on-premise S3-хранилищу часто возникают ошибки. Большинство из них имеют конкретную причину, которую можно найти по типу исключения.
HTTP 405 Method Not Allowed - самая коварная ошибка. Она возникает, когда Spark пытается использовать API-метод, который данное хранилище не реализовало или реализовало иначе. Типичные причины:
- Включён Magic Committer на хранилище без поддержки magic files API (Ceph, SeaweedFS)
- Хранилище не поддерживает
DELETEнескольких объектов за один запрос (DeleteObjects) - Хранилище не поддерживает
CopyObject(используется Directory Committer для rename)
# Диагностический скрипт: проверяем базовые S3 операции
import boto3
from botocore.config import Config
import time
def diagnose_s3_compatibility(endpoint_url: str, access_key: str, secret_key: str) -> dict:
"""
Проверяет совместимость хранилища с операциями, которые использует Spark.
Запускайте перед первым развёртыванием пайплайна на новом хранилище.
"""
client = boto3.client(
"s3",
endpoint_url=endpoint_url,
aws_access_key_id=access_key,
aws_secret_access_key=secret_key,
config=Config(signature_version="s3v4"),
)
results = {}
test_bucket = f"spark-compat-test-{int(time.time())}"
try:
# 1. CREATE BUCKET
client.create_bucket(Bucket=test_bucket)
results["create_bucket"] = "✅ OK"
# 2. PUT OBJECT (имитирует запись Parquet-файла)
test_data = b"x" * (5 * 1024 * 1024) # 5 MB
client.put_object(Bucket=test_bucket, Key="test/part-0.parquet", Body=test_data)
results["put_object"] = "✅ OK"
# 3. HEAD OBJECT (S3A делает HEAD перед каждым GET для проверки размера)
client.head_object(Bucket=test_bucket, Key="test/part-0.parquet")
results["head_object"] = "✅ OK"
# 4. LIST OBJECTS V2 (используется при FileSystem.listStatus)
resp = client.list_objects_v2(Bucket=test_bucket, Prefix="test/")
assert len(resp.get("Contents", [])) > 0
results["list_objects_v2"] = "✅ OK"
# 5. GET OBJECT с Range (range request - ключевой для Parquet column pruning)
client.get_object(Bucket=test_bucket, Key="test/part-0.parquet",
Range="bytes=0-1023")
results["get_object_range"] = "✅ OK"
# 6. MULTIPART UPLOAD (используется для файлов > multipart.size)
mpu = client.create_multipart_upload(Bucket=test_bucket, Key="test/large.parquet")
upload_id = mpu["UploadId"]
part = client.upload_part(
Bucket=test_bucket, Key="test/large.parquet",
UploadId=upload_id, PartNumber=1,
Body=test_data
)
client.complete_multipart_upload(
Bucket=test_bucket, Key="test/large.parquet",
UploadId=upload_id,
MultipartUpload={"Parts": [{"PartNumber": 1, "ETag": part["ETag"]}]}
)
results["multipart_upload"] = "✅ OK"
# 7. DELETE OBJECTS (batch delete - Directory Committer удаляет staging)
client.delete_objects(
Bucket=test_bucket,
Delete={"Objects": [
{"Key": "test/part-0.parquet"},
{"Key": "test/large.parquet"},
]}
)
results["delete_objects_batch"] = "✅ OK"
# 8. COPY OBJECT (используется Directory Committer для rename)
client.put_object(Bucket=test_bucket, Key="src/file.parquet", Body=b"data")
client.copy_object(
CopySource={"Bucket": test_bucket, "Key": "src/file.parquet"},
Bucket=test_bucket,
Key="dst/file.parquet"
)
results["copy_object"] = "✅ OK"
except Exception as e:
results[f"FAILED"] = f"❌ {type(e).__name__}: {e}"
finally:
# Очистка тестового бакета
try:
paginator = client.get_paginator("list_objects_v2")
for page in paginator.paginate(Bucket=test_bucket):
objects = page.get("Contents", [])
if objects:
client.delete_objects(
Bucket=test_bucket,
Delete={"Objects": [{"Key": o["Key"]} for o in objects]}
)
client.delete_bucket(Bucket=test_bucket)
except Exception:
pass
return results
# Пример использования:
# results = diagnose_s3_compatibility(
# "http://minio.example.com:9000", "minioadmin", "minioadmin"
# )
# for op, status in results.items():
# print(f"{op:30s}: {status}")
Лабораторное задание: MinIO vs SeaweedFS бенчмарк¶
Следующий скрипт запускает один и тот же Spark-пайплайн на двух хранилищах и сравнивает время выполнения. Он предназначен для запуска в Docker Compose окружении с поднятыми инстансами MinIO и SeaweedFS.
# docker-compose.yml для лабораторной работы
services:
minio:
image: minio/minio:latest
command: server /data --console-address :9001
environment:
MINIO_ROOT_USER: minioadmin
MINIO_ROOT_PASSWORD: minioadmin
ports:
- "9000:9000"
- "9001:9001"
volumes:
- minio_data:/data
seaweedfs-master:
image: chrislusf/seaweedfs:latest
command: master -mdir=/data
ports:
- "9333:9333"
volumes:
- seaweed_master:/data
seaweedfs-volume:
image: chrislusf/seaweedfs:latest
command: volume -mserver=seaweedfs-master:9333 -dir=/data -port=8080
volumes:
- seaweed_volume:/data
depends_on:
- seaweedfs-master
seaweedfs-filer:
image: chrislusf/seaweedfs:latest
command: filer -master=seaweedfs-master:9333 -s3
ports:
- "8333:8333" # S3 API
depends_on:
- seaweedfs-master
- seaweedfs-volume
volumes:
minio_data:
seaweed_master:
seaweed_volume:
# benchmark_storages.py
import time
import random
import boto3
from botocore.config import Config
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType
STORAGES = {
"minio": {
"endpoint": "http://localhost:9000",
"access_key": "minioadmin",
"secret_key": "minioadmin",
"magic_committer": True,
},
"seaweedfs": {
"endpoint": "http://localhost:8333",
"access_key": "any",
"secret_key": "any",
"magic_committer": False,
},
}
def generate_test_data(spark: SparkSession, n_rows: int):
"""Генерирует синтетические данные для бенчмарка."""
event_types = ["purchase", "page_view", "search", "add_to_cart", "checkout"]
countries = ["RU", "DE", "US", "FR", "BR", "CN"]
return spark.range(n_rows).select(
F.concat(F.lit("evt-"), F.col("id").cast("string")).alias("event_id"),
(F.rand() * 100000).cast("long").alias("user_id"),
F.element_at(
F.array([F.lit(e) for e in event_types]),
(F.rand() * len(event_types) + 1).cast("int")
).alias("event_type"),
(F.rand() * 500).alias("amount"),
F.element_at(
F.array([F.lit(c) for c in countries]),
(F.rand() * len(countries) + 1).cast("int")
).alias("country"),
F.current_timestamp().alias("event_time"),
)
def run_benchmark(storage_name: str, config: dict, n_rows: int = 2_000_000) -> dict:
"""
Запускает полный ETL пайплайн на указанном хранилище и измеряет время каждого этапа.
Pipeline: генерация данных → запись CSV → чтение → агрегация → запись Parquet
"""
print(f"\n{'='*60}")
print(f"Бенчмарк: {storage_name.upper()} ({config['endpoint']})")
print(f"Строк: {n_rows:,}")
# Создаём boto3 клиент для управления бакетами
s3 = boto3.client(
"s3",
endpoint_url=config["endpoint"],
aws_access_key_id=config["access_key"],
aws_secret_access_key=config["secret_key"],
config=Config(signature_version="s3v4"),
)
# Создаём тестовые бакеты
for bucket in ["bench-raw", "bench-silver"]:
try:
s3.create_bucket(Bucket=bucket)
except Exception:
pass # бакет уже существует
# Настраиваем SparkSession
committer = "magic" if config["magic_committer"] else "directory"
spark = (
SparkSession.builder
.master("local[4]")
.appName(f"benchmark-{storage_name}")
.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", config["endpoint"])
.config("spark.hadoop.fs.s3a.path.style.access", "true")
.config("spark.hadoop.fs.s3a.aws.credentials.provider",
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider")
.config("spark.hadoop.fs.s3a.access.key", config["access_key"])
.config("spark.hadoop.fs.s3a.secret.key", config["secret_key"])
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false")
.config("spark.hadoop.fs.s3a.committer.name", committer)
.config("spark.hadoop.fs.s3a.committer.magic.enabled",
"true" if committer == "magic" else "false")
.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.sql.shuffle.partitions", "8")
.config("spark.ui.enabled", "false")
.getOrCreate()
)
spark.sparkContext.setLogLevel("ERROR")
timings = {}
# ── Этап 1: генерация и запись JSON (имитация Bronze ingestion) ──────────
raw_path = f"s3a://bench-raw/events/"
t0 = time.time()
df = generate_test_data(spark, n_rows)
df.repartition(20).write.mode("overwrite").json(raw_path)
timings["write_json"] = time.time() - t0
print(f" Запись {n_rows:,} строк JSON: {timings['write_json']:.1f}s")
# ── Этап 2: чтение JSON, трансформация, запись Parquet ──────────────────
silver_path = f"s3a://bench-silver/events/"
t0 = time.time()
raw_df = spark.read.json(raw_path)
silver_df = (
raw_df
.dropDuplicates(["event_id"])
.filter(F.col("event_type").isNotNull())
.withColumn("event_type", F.lower(F.trim(F.col("event_type"))))
.repartition(8)
)
silver_df.write.mode("overwrite").partitionBy("country").parquet(silver_path)
timings["etl_json_to_parquet"] = time.time() - t0
print(f" ETL JSON → Parquet: {timings['etl_json_to_parquet']:.1f}s")
# ── Этап 3: агрегация из Parquet (имитация Gold) ─────────────────────────
t0 = time.time()
result_df = (
spark.read.parquet(silver_path)
.groupBy("country", "event_type")
.agg(
F.count("*").alias("event_count"),
F.countDistinct("user_id").alias("unique_users"),
F.sum("amount").alias("total_amount"),
)
)
row_count = result_df.count()
timings["aggregation"] = time.time() - t0
print(f" Агрегация ({row_count} групп): {timings['aggregation']:.1f}s")
# ── Этап 4: LIST директории (имитация planning phase) ────────────────────
t0 = time.time()
files = spark.read.parquet(silver_path).inputFiles()
timings["list_files"] = time.time() - t0
print(f" LIST {len(files)} файлов: {timings['list_files']:.3f}s")
timings["total"] = sum(v for v in timings.values())
print(f" ИТОГО: {timings['total']:.1f}s")
spark.stop()
return timings
def print_comparison(results: dict) -> None:
"""Выводит сравнительную таблицу результатов."""
print(f"\n{'='*60}")
print("СРАВНЕНИЕ ХРАНИЛИЩ")
print(f"{'Этап':<30} {'MinIO':>12} {'SeaweedFS':>12} {'Разница':>12}")
print("-" * 70)
minio_t = results.get("minio", {})
seaweed_t = results.get("seaweedfs", {})
stages = ["write_json", "etl_json_to_parquet", "aggregation", "list_files", "total"]
stage_labels = {
"write_json": "Запись JSON",
"etl_json_to_parquet": "ETL JSON→Parquet",
"aggregation": "Gold агрегация",
"list_files": "LIST файлов",
"total": "ИТОГО",
}
for stage in stages:
minio_val = minio_t.get(stage, 0)
seaweed_val = seaweed_t.get(stage, 0)
if minio_val > 0:
diff_pct = (seaweed_val - minio_val) / minio_val * 100
diff_str = f"{diff_pct:+.0f}%"
else:
diff_str = "N/A"
label = stage_labels[stage]
print(f"{' ' if stage != 'total' else ''}{label:<28} "
f"{minio_val:>11.1f}s {seaweed_val:>11.1f}s {diff_str:>12}")
if __name__ == "__main__":
N_ROWS = 1_000_000 # 1 миллион событий
results = {}
for name, cfg in STORAGES.items():
try:
results[name] = run_benchmark(name, cfg, N_ROWS)
except Exception as e:
print(f"Ошибка на {name}: {e}")
if len(results) == 2:
print_comparison(results)
Интерпретация результатов бенчмарка. Ожидаемые паттерны:
- Запись JSON (write_json): MinIO и SeaweedFS покажут близкие результаты для файлов ~10-50 MB каждый. Если SeaweedFS заметно медленнее - проблема с конфигурацией Filer или network.
- ETL JSON→Parquet: здесь проявляется разница Magic Committer (MinIO) vs Directory Committer (SeaweedFS). MinIO должен быть на 15-30% быстрее на стадии commit.
- LIST файлов: для 20 файлов разница незаметна. Увеличьте
repartition(200)для честного теста LIST производительности. - Агрегация: почти не зависит от хранилища - это чистые Spark CPU операции после загрузки данных в память Executor'ов.
Смотрите в Spark UI вкладку Storage → SQL → конкретный запрос → Metrics для каждого Stage. Там видно время каждой задачи: task deserialization time (чтение с S3), compute time (CPU), write time (запись в S3).
Emerging alternatives: Versity Gateway и RustFS¶
Ландшафт S3-совместимых хранилищ не ограничивается пятью системами, рассмотренными выше. Два относительно новых проекта заслуживают отдельного упоминания - не как полноценная замена MinIO для Data Lake, но как решения для специфических задач.
Versity S3 Gateway: POSIX-хранилище с S3-интерфейсом¶
Versity Gateway - это S3-шлюз, который работает поверх стандартных POSIX файловых систем (ext4, XFS, ZFS, btrfs). В отличие от MinIO или Ceph, которые управляют данными самостоятельно, Versity просто предоставляет S3-интерфейс к уже существующей файловой системе. Метаданные объектов хранятся в расширенных атрибутах файлов (xattr), а содержимое объекта - как обычный файл.
Versity Gateway (s3gw process)
↓
Обычная директория на XFS / ZFS / ext4
↓
Физический диск или NFS-примонтированное хранилище
Где это применяется. Versity используется в национальных лабораториях США (Sandia National Labs, Los Alamos) и академических учреждениях, где уже существует большая POSIX-инфраструктура - HPC кластеры с Lustre, NFS-массивы - и нужно добавить S3 API без полной замены storage backend. По заявленным характеристикам, Versity достигает скоростей, близких к пропускной способности сети, на LAN.
Для Data Engineering. Versity не является решением для построения нового Data Lakehouse с нуля. Но если у вашей организации есть legacy NAS-массив или HPC-хранилище с доступом по POSIX, и вы хотите прочитать данные с него через стандартный Spark S3A Connector без переноса данных - Versity может быть мостом.
Ограничения: производительность напрямую зависит от производительности underlying файловой системы; нет встроенного Erasure Coding; масштабирование ограничено возможностями POSIX-слоя.
RustFS: прямой бинарный преемник MinIO¶
RustFS - реализация S3-совместимого объектного хранилища на Rust с архитектурой, близкой к MinIO. Позиционируется как drop-in replacement с тем же форматом конфигурации и теми же путями API. Ключевые характеристики:
- Написан на Rust - меньший memory footprint, отсутствие GC пауз
- Прямая работа с блочными устройствами или директориями файловой системы
- Многодисковое распределение данных с автоматическим восстановлением
Практически RustFS пока не достиг уровня зрелости MinIO: нет поддержки FreeBSD (только Linux), документация минимальна, production-кейсов немного. Производительность на крупных объектах сравнима с MinIO, но на мелких объектах (файловые операции через Rust syscall) есть деградация по сравнению с оптимизированным Go-кодом MinIO.
Если MinIO по каким-то причинам не подходит (лицензионные ограничения, требование полностью Rust-стека в инфраструктуре) - RustFS стоит рассмотреть как запасной вариант.
Итоговые рекомендации по выбору¶
Выбор хранилища определяется тремя осями: масштаб данных, операционные возможности команды и специфика нагрузки.
Для нового аналитического проекта с объёмом данных до сотен терабайт, без специализированных инженеров хранилищ и стандартной нагрузкой (крупные Parquet файлы, Medallion Architecture) - MinIO остаётся безусловным первым выбором. Он обеспечивает наилучшую интеграцию со Spark, поддерживает Magic Committer, прост в эксплуатации.
Ceph RGW уместен, когда организация уже управляет Ceph кластером для виртуальных машин или когда нужно масштабироваться в петабайты с опытной командой.
SeaweedFS даёт реальное преимущество именно в сценарии Bronze-слоя с высокочастотным стримингом мелких файлов. В комбинации с MinIO для Silver/Gold - интересная гибридная архитектура.
Apache Ozone - выбор Hadoop-организаций, которые упёрлись в ограничение NameNode и хотят сохранить совместимость с Hive, Trino, YARN.
Garage остаётся нишевым инструментом для геораспределённого хранения - он не подходит для single-datacenter аналитических нагрузок из-за eventual consistency и отсутствия Magic Committer.