HDFS vs S3: когда что выбирать - latency, consistency, стоимость
Архитектурное сравнение HDFS и S3-совместимых хранилищ для Data Engineer: философия Co-located vs Decoupled, битва latency и throughput, гарантии consistency и атомарность rename, влияние на Spark Committers и Lakehouse форматы, модель TCO, матрица выбора и три кейс-стади из enterprise-практики.
1. Архитектурная парадигма: Co-located vs Decoupled Storage¶
Прежде чем сравнивать конкретные технические характеристики HDFS и S3, необходимо понять фундаментальную философию, лежащую в основе каждой системы. Это не просто два способа хранить файлы - это два принципиально разных взгляда на то, как должны взаимодействовать вычисления и данные.
Философия HDFS: принести код к данным¶
HDFS создавался в эпоху, когда сеть была дорогим и медленным ресурсом. В 2006–2008 годах типичная скорость сетевого соединения в дата-центре составляла 1 Gbps - это ~125 MB/s. Диски HDD давали 50–100 MB/s последовательного чтения. Передавать 1 TB данных по такой сети - 2+ часа.
Hadoop предложил революционное решение: не перемещай данные к вычислениям - перемещай вычисления к данным. Каждый сервер кластера одновременно является и хранилищем (DataNode) и вычислительным узлом (NodeManager). Задачи MapReduce или Spark Tasks запускаются на тех же серверах, где физически лежат нужные блоки данных. Сетевой трафик сводится к минимуму.
Философия S3: бесконечное хранилище, эфемерные вычисления¶
Объектное хранилище строится на противоположной идее: compute и storage - это независимые ресурсы, которые должны масштабироваться раздельно. Вы арендуете вычисления на время обработки (Spot Instances, Kubernetes pods), а данные хранятся в отдельном объектном хранилище - постоянно и независимо от наличия вычислительных ресурсов.
Когда Spark Job завершилась - кластер можно полностью остановить (Scale-to-Zero). Данные в S3 остаются. Это радикально меняет экономику: вы не платите за «горячий» кластер HDFS в период, когда данные просто лежат и ждут.
Смена вех: как 100 Gbps сети изменили правила игры¶
Главное историческое преимущество HDFS - устранение сетевого bottleneck за счёт локальности - стало менее значимым с распространением высокоскоростных сетей.
| Эпоха | Сеть дата-центра | HDFS/NVMe диск | Преимущество HDFS |
|---|---|---|---|
| 2006–2012 | 1 GbE (125 MB/s) | HDD ~80 MB/s | Огромное (NVMe > Сеть) |
| 2012–2018 | 10 GbE (1.25 GB/s) | SSD ~500 MB/s | Значительное |
| 2018–2022 | 25 GbE (3 GB/s) | NVMe ~3 GB/s | Небольшое |
| 2022–сейчас | 100 GbE (12.5 GB/s) + RDMA | NVMe ~7 GB/s | Минимальное для seq. read |
При 100 GbE сеть перестала быть узким местом для sequential чтения больших Parquet-файлов. Spark может читать данные из S3 почти так же быстро, как с локального NVMe. Это объясняет, почему S3-based Lakehouse стал индустриальным стандартом: физическое преимущество HDFS нивелировалось прогрессом сетевого оборудования.
2. Битва Latency: HTTP REST API против RPC в оперативную память¶
Несмотря на рост пропускной способности сетей, в области latency (задержки на отдельные операции) HDFS по-прежнему имеет существенное преимущество. И это преимущество критично для metadata-интенсивных нагрузок.
Тайминги metadata-операций¶
Операции с metadata (LIST директорий, HEAD объектов, CREATE/DELETE файлов) - это критический путь планирования Spark Job. Перед тем как читать данные, Spark Driver должен перечислить файлы и получить их блочные адреса.
При 10 000 файлов в партиции:
- HDFS: 10 LIST запросов × 2 мс = 20 мс на обнаружение файлов
- S3: 10 LIST запросов × 50 мс = 500 мс на обнаружение файлов
Разница в 25 раз только на metadata-фазе. При миллионе файлов разница становится катастрофической: 2 секунды на HDFS vs 50 секунд на S3.
S3 Throttling: неожиданный bottleneck при высоком параллелизме¶
AWS S3 (и большинство S3-совместимых хранилищ) применяют rate limiting по prefix:
- GET/HEAD: 5500 request/s на prefix
- PUT/DELETE: 3500 request/s на prefix
При запуске 200 Spark Executor'ов, каждый из которых выполняет 50 metadata-запросов в секунду, возникает 10 000 req/s - в два раза больше лимита. Результат: HTTP 503 SlowDown ошибки и exponential backoff в S3A Connector.
# Признак S3 throttling в логах Spark:
# WARN RetryThrottlingHandler: Throttling exception from S3. Sleeping for...
# WARN AmazonS3Exception: Status Code: 503, AWS Service: Amazon S3,
# AWS Request ID: ..., Extended Request ID: ...,
# Error Message: Slow Down, Error type: Service
# Решение: S3 автоматически увеличивает лимит при хорошей партиционной схеме
# Ключ: распределять объекты по разным prefix'ам
# НЕ делать: s3://bucket/data/events/ (все объекты в одном prefix)
# ДЕЛАТЬ: s3://bucket/data/events/year=2024/month=01/day=15/ (разные prefix)
Максимальный throughput: Short-Circuit Local Reads vs S3 Streaming¶
При идеальной конфигурации (NVMe диски + Short-Circuit) HDFS может давать пропускную способность, недостижимую для S3:
- HDFS Short-Circuit (NVMe): 3-7 GB/s на один Executor, данные читаются напрямую с локального NVMe через Unix Domain Socket, минуя сетевой стек полностью
- S3 через 100 GbE: 1-2 GB/s на один Executor, данные всегда идут по сети
Для итеративных ML-рабочих нагрузок (обучение модели за 100 эпох, каждая эпоха - полное сканирование датасета) и heavy Shuffle (запись временных файлов на диск при sort-merge join на сотнях GB) HDFS остаётся быстрее за счёт локальных NVMe дисков.
3. Гарантии согласованности и проблема атомарности¶
Consistency (согласованность) - это одно из самых важных и наименее очевидных различий между HDFS и S3. Неправильное понимание consistency может приводить к тихим ошибкам в данных, которые обнаруживаются только спустя дни или недели.
Строгая консистентность HDFS¶
HDFS предоставляет строгую read-after-write consistency по дизайну. После того как запись завершена (write() + close()), файл немедленно виден всем клиентам - в том числе работающим на других серверах кластера.
Это следует из централизованной архитектуры: все знания о файловой системе хранятся в одном месте (NameNode), и все клиенты обращаются к одному источнику правды. Нет ни кешей, ни репликации метаданных, которые могли бы создать рассогласование.
Историческая боль eventual consistency в S3¶
Исходный дизайн AWS S3 (2006–2020) использовал eventual consistency для определённых операций:
PUTнового объекта → strong consistency (объект сразу виден)PUTUPDATE существующего объекта → eventual consistency (обновление может быть не видно немедленно)DELETEобъекта → eventual consistency (старый объект может ещё возвращаться в LIST)LISTпосле bulk PUT → eventual consistency (не все объекты могут появиться в LIST)
Это создавало серьёзные проблемы для Spark. Классический сценарий:
Spark Job 1: записывает 1000 файлов в s3://bucket/output/
Spark Job 1: завершается успешно, пишет _SUCCESS файл
Spark Job 2: немедленно запускается, делает LIST s3://bucket/output/
Spark Job 2: видит только 700 из 1000 файлов (еще 300 не "проявились")
Spark Job 2: обрабатывает неполные данные → тихая ошибка в данных
В декабре 2020 AWS ввёл strong consistency для всех операций S3. Это стало важным поворотным моментом. Теперь AWS S3 предоставляет read-after-write consistency для PUT, LIST и DELETE.
Проблема атомарности: главное фундаментальное различие¶
Это наиболее принципиальное техническое различие между HDFS и S3, которое влияет на все операции финализации записи в Spark.
В HDFS: папки - это объекты первого класса в namespace NameNode. rename("/tmp/output_tmp", "/data/output") - это изменение одной строки в метаданных NameNode. Операция занимает O(1) по времени независимо от числа файлов внутри папки.
В S3: папок не существует. S3 - это плоское пространство ключей. s3://bucket/tmp/output_tmp/file.parquet - это просто объект с ключом tmp/output_tmp/file.parquet. «Переименование папки» в S3 - это:
- GET list всех объектов с префиксом
tmp/output_tmp/ - PUT каждый объект под новым ключом
data/output/... - DELETE каждый старый объект
Это O(N) операция где N - число файлов. Для директории с 10 000 файлами это 20 000+ API-запросов, которые занимают секунды или минуты.
# Последствия неатомарного rename в S3 для Spark:
# Стандартный FileOutputCommitter (v1) работает так:
# 1. Spark Tasks записывают в s3://bucket/_temporary/attempt_xxx/
# 2. Driver делает rename/_temporary/ → s3://bucket/output/
# В S3 шаг 2 = тысячи PUT + тысячи DELETE операций
# Если кластер упадёт в середине этой операции:
# - Часть файлов в новом пути (s3://bucket/output/)
# - Часть файлов в старом пути (s3://bucket/_temporary/)
# - Данные неконсистентны! Нет транзакционности!
# Решение: S3A Magic Committer или Staging Committer
# (подробнее в модуле 05 о S3 хранилищах)
from pyspark.sql import SparkSession
spark = SparkSession.builder \
# Magic Committer: файлы сразу пишутся по финальному пути
# Нет staging директории = нет rename-проблемы!
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.fs.s3a.committer.magic.enabled", "true") \
.config("spark.sql.sources.commitProtocolClass",
"org.apache.spark.internal.io.cloud.PathOutputCommitProtocol") \
.getOrCreate()
4. Влияние на жизненный цикл Spark Job: Committers и Lakehouse форматы¶
Выбор хранилища влияет не только на производительность чтения, но и на всю логику записи данных и планирования запросов в Spark.
Как хранилище меняет физический план Catalyst¶
Catalyst Optimizer в Spark генерирует физический план запроса, учитывая характеристики источника данных. Ключевое различие:
Практическое следствие для настройки: при работе с S3 параметр spark.locality.wait нужно выставлять в 0. Ожидание Node-Local слота бессмысленно - такого слота никогда не будет. Без этой настройки Spark будет тратить несколько секунд на каждый Stage, ожидая несуществующую локальность.
spark = SparkSession.builder \
# S3: никогда нет NODE_LOCAL → ожидание бессмысленно
.config("spark.locality.wait", "0") \
.config("spark.locality.wait.node", "0") \
.config("spark.locality.wait.process", "0") \
# HDFS: ждём NODE_LOCAL слот
# .config("spark.locality.wait.node", "10s")
.getOrCreate()
Разбор Committers: почему стандартный не работает с S3¶
FileOutputCommitter - это механизм финализации записи в Spark. Он гарантирует атомарность: либо все данные записаны корректно, либо ни одного файла.
Для HDFS FileOutputCommitter v1 - идеальный выбор. Rename директории атомарен и мгновенен.
Для S3 варианты:
- Magic Committer: объекты записываются напрямую в финальный путь через специальный механизм S3. Нет staging → нет rename → нет проблемы. Требует поддержки хранилища (MinIO, AWS S3).
- Staging Committer: данные записываются на локальный диск Executor'а, затем копируются напрямую в финальный путь S3. Нет staging директории в S3.
- Directory Committer: компромисс - staging в S3, но rename на уровне Directory (не файлов).
Lakehouse как великий уравнитель¶
Apache Iceberg, Delta Lake и Apache Hudi кардинально меняют уравнение HDFS vs S3. Эти форматы решают проблемы S3, которые исторически делали его неудобным для аналитики:
Проблема 1 - дорогой LIST: Без Lakehouse-форматов Spark делает LIST всей директории при каждом чтении. При миллионе файлов - тысячи LIST запросов к S3.
Iceberg решение: Manifest List файл содержит точный список всех Parquet-файлов. Spark читает один manifest → получает все пути → читает данные напрямую. Нет LIST запросов.
Проблема 2 - неатомарная запись: Без Lakehouse файлы появляются частично при concurrent writes.
Delta Lake решение: _delta_log/ содержит атомарные commit записи. Читатель видит только «полные» состояния таблицы через snapshot isolation.
# Iceberg с S3: никаких LIST запросов при чтении!
spark = SparkSession.builder \
.config("spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
.config("spark.sql.catalog.iceberg",
"org.apache.iceberg.spark.SparkCatalog") \
.config("spark.sql.catalog.iceberg.type", "hive") \
.config("spark.sql.catalog.iceberg.warehouse", "s3a://warehouse/") \
# S3FileIO: обходит Hadoop S3A, работает напрямую через AWS SDK
.config("spark.sql.catalog.iceberg.io-impl",
"org.apache.iceberg.aws.s3.S3FileIO") \
.config("spark.sql.catalog.iceberg.s3.endpoint", "http://minio:9000") \
.getOrCreate()
# Spark читает metadata.json → manifest list → manifest files
# Ни одного LIST запроса к S3!
df = spark.table("iceberg.warehouse.events")
df.filter("event_date = '2024-01-15'").explain()
# Physical Plan: BatchScan iceberg[event_date=2024-01-15] (точный список файлов)
# Нет: ListObjectsV2 → FilteredScan → ...
5. Экономика данных: детальная модель TCO¶
Total Cost of Ownership (TCO) - часто решающий фактор при выборе архитектуры. Рассмотрим реалистичную стоимостную модель для типичного Enterprise Data Lake.
CapEx HDFS: видимые и скрытые затраты¶
Стартовые капитальные затраты для HDFS-кластера на 1 PB полезных данных (replication 3x → 3 PB raw):
| Компонент | Конфигурация | Стоимость (approx.) |
|---|---|---|
| DataNode серверы | 10 × 40 TB NVMe (400 TB raw) × 3 стойки | ~$150K |
| NameNode серверы | 2 × High RAM (256 GB) | ~$30K |
| JournalNode + ZK | 3 × средних сервера | ~$15K |
| Сетевое оборудование | 100 GbE + коммутаторы | ~$50K |
| Итого hardware | ~$245K | |
| Инсталляция + настройка | ~$30K | |
| Итого CapEx | ~$275K |
Это только первоначальные затраты. Дополнительно ежегодно:
- Электричество (40W × 15 серверов × 24h × 365 days): ~$20K/год
- Инженеры Hadoop Ops (1-2 человека): ~$150K–$300K/год
- Замена дисков (~10% в год): ~$20K/год
- Итого OpEx: ~$200–350K/год
OpEx S3: Pay-as-you-go и скрытые API-расходы¶
Стоимость хранения 1 PB в AWS S3 Standard:
1 PB × $0.023/GB/месяц = $23,552/месяц = $282,624/год
Для cold данных (S3 Glacier Instant Retrieval):
1 PB × $0.004/GB/месяц = $4,096/месяц = $49,152/год
Скрытые расходы на S3 API - часто неожиданны для команд, не ведущих мониторинг:
# Расчёт стоимости S3 API операций для типичного Spark ETL
# Допущения:
# - 100 Spark Job в день
# - Каждый Job читает 10,000 файлов
# - 5 LIST запросов на директорию (paginated)
# - replication: нет (один экземпляр данных в S3)
daily_jobs = 100
files_per_job = 10_000
list_requests_per_job = files_per_job // 1000 * 5 # S3 returns 1000 obj/request
# Стоимость операций (AWS S3, us-east-1, 2024):
s3_get_price_per_1000 = 0.0004 # $0.0004 за 1000 GET запросов
s3_list_price_per_1000 = 0.005 # $0.005 за 1000 LIST запросов (в 12.5 раз дороже!)
daily_get_requests = daily_jobs * files_per_job
daily_list_requests = daily_jobs * list_requests_per_job
daily_api_cost = (
daily_get_requests / 1000 * s3_get_price_per_1000 +
daily_list_requests / 1000 * s3_list_price_per_1000
)
print(f"GET запросов в день: {daily_get_requests:,}")
print(f"LIST запросов в день: {daily_list_requests:,}")
print(f"Стоимость API в день: ${daily_api_cost:.2f}")
print(f"Стоимость API в год: ${daily_api_cost * 365:.2f}")
# Пример вывода:
# GET запросов в день: 1,000,000
# LIST запросов в день: 5,000
# Стоимость API в день: $0.43
# Стоимость API в год: $156.95
# Это разумно. НО при Small Files (миллион мелких файлов):
daily_jobs_bad = 100
files_per_job_bad = 1_000_000 # Small Files!
list_requests_bad = files_per_job_bad // 1000 * 5
daily_get_bad = daily_jobs_bad * files_per_job_bad
daily_list_bad = daily_jobs_bad * list_requests_bad
daily_api_cost_bad = (
daily_get_bad / 1000 * s3_get_price_per_1000 +
daily_list_bad / 1000 * s3_list_price_per_1000
)
print(f"\nПри Small Files (1M файлов на Job):")
print(f"Стоимость API в день: ${daily_api_cost_bad:.2f}")
print(f"Стоимость API в год: ${daily_api_cost_bad * 365:,.0f}")
# При Small Files:
# Стоимость API в день: $42.50
# Стоимость API в год: $15,512 ← просто за мелкие файлы!
Сравнительная TCO-модель для 1 PB за 3 года¶
На уровне 1 PB HDFS и S3 сопоставимы по TCO при грамотном управлении. Основное экономическое преимущество S3 проявляется в:
- Scale-to-Zero: вычислительный кластер существует только во время обработки. HDFS требует постоянно работающих DataNode серверов.
- Cold Storage tiering: S3 Glacier (3-5x дешевле S3 Standard) для архивных данных. HDFS не имеет встроенного tiering.
- Отсутствие Hadoop Ops команды: экономия $150-300K/год на специализированных инженерах.
6. Матрица выбора: сводный гид для Data-архитектора¶
Полная сравнительная таблица¶
| Критерий | HDFS (On-Premise) | S3 / Object Storage |
|---|---|---|
| Data Locality | ✅ NODE_LOCAL, PROCESS_LOCAL | ❌ Всегда ANY |
| Metadata latency | ✅ 0.1-2 мс (RAM) | ⚠️ 20-100 мс (HTTP) |
| Sequential throughput | ✅ 500-7000 MB/s (NVMe+SC) | ⚠️ 100-2000 MB/s (сеть) |
| Rename атомарность | ✅ O(1) мгновенно | ❌ O(N) COPY+DELETE |
| Strong Consistency | ✅ Всегда | ✅ AWS S3 (c 2020), MinIO |
| Масштабируемость | ⚠️ Ограничена NameNode RAM | ✅ Практически безлимитная |
| Scale-to-Zero Compute | ❌ DataNode должны работать | ✅ Storage независим |
| Cold Tiering | ❌ Только EC (50% экономии) | ✅ Glacier (80% экономии) |
| Stоимость хранения | ❌ 3x replication overhead | ✅ Платишь за реальный объём |
| Операционная сложность | ❌ NameNode, JN, ZK, Balancer | ✅ Managed service |
| Spark Committer | ✅ FileOutputCommitter | ⚠️ Нужен Magic/Staging Committer |
| Lakehouse (Iceberg) | ✅ Отлично работает | ✅ Спроектировано для S3 |
| Short-Circuit Reads | ✅ Доступен | ❌ Невозможен |
| Регуляторный контроль | ✅ Полный контроль on-prem | ⚠️ Зависит от провайдера |
| Time-to-Market | ❌ Недели/месяцы настройки | ✅ Часы (managed cloud) |
Чек-лист принятия архитектурного решения¶
ВЫБИРАЙТЕ HDFS если:
□ Уже есть зрелый Hadoop-кластер с опытной командой Ops
□ Данные не могут покидать корпоративный периметр (банки, госсектор)
□ Нагрузка - тяжёлый Shuffle или итеративные ML-алгоритмы
□ Данные используются 24/7 без простоев (нет смысла в Scale-to-Zero)
□ Требуется POSIX-семантика (rename, hardlinks, директории)
□ Compliance/суверенитет данных: полный физический контроль
□ Объём данных: 1-100 PB, стабильный рост
ВЫБИРАЙТЕ S3/Object Storage если:
□ Новый проект без legacy Hadoop-инфраструктуры
□ Нужен Scale-to-Zero (вычисления нужны только N часов в день)
□ Данные неизменяемы (write-once, read-many)
□ Используется Iceberg/Delta Lake (нивелируют недостатки S3)
□ Важна скорость развёртывания (облачный managed service)
□ Объём хранения непредсказуем или растёт быстро
□ Cold данные занимают > 50% общего объёма (Glacier tiering)
□ Команда небольшая, нет ресурсов на Hadoop Ops
7. Практика: три кейс-стади из Enterprise¶
Кейс A: Телеком - архив логов сетевого оборудования¶
Описание: крупный телеком-оператор собирает логи с 50 000 единиц сетевого оборудования. Объём: 5 PB в месяц, хранение 3 года. Доступ - раз в квартал для регулятора. Существует требование хранить данные в России.
Анализ:
# Параметры:
monthly_volume_tb = 5_000 # 5 PB/месяц
retention_years = 3
access_frequency = "quarterly" # редко
data_sovereignty = "Russia" # no-cloud
# Подсчёт объёма:
total_tb = monthly_volume_tb * 12 * retention_years
print(f"Итого данных: {total_tb} TB = {total_tb/1000:.0f} PB")
# Итого данных: 180000 TB = 180 PB
# HDFS:
hdfs_overhead = 3.0 # 3x репликация
hdfs_raw_tb = total_tb * hdfs_overhead
print(f"HDFS raw storage: {hdfs_raw_tb/1000:.0f} PB (с 3x репл.)")
# MinIO с EC RS-10-4:
ec_overhead = 1.4 # 40% overhead
minio_raw_tb = total_tb * ec_overhead
print(f"MinIO EC RS-10-4: {minio_raw_tb/1000:.0f} PB")
# Экономия: 180 PB vs 252 PB = 72 PB экономии!
Рекомендация: MinIO с Erasure Coding RS-10-4 поверх Iceberg-таблиц.
- Суверенитет данных: on-premise, в России ✅
- Стоимость: 40% overhead vs 200% у HDFS репликации ✅
- Холодные данные (доступ раз в квартал): EC с высоким соотношением данных к parity ✅
- Iceberg: устраняет проблему листинга миллиардов log-файлов ✅
- Spark Committer: Magic Committer для MinIO ✅
Конфигурация Spark для этого кейса:
spark = SparkSession.builder \
.appName("telecom-compliance-query") \
# MinIO endpoint
.config("spark.hadoop.fs.s3a.endpoint", "http://minio.internal:9000") \
.config("spark.hadoop.fs.s3a.path.style.access", "true") \
# Magic Committer для MinIO
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.fs.s3a.committer.magic.enabled", "true") \
# Iceberg каталог
.config("spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
# Для compliance запросов: большие executors, долгие таймауты
.config("spark.executor.memory", "32g") \
.config("spark.sql.files.maxPartitionBytes", str(512 * 1024 * 1024)) \
# S3 locality нет → не ждём
.config("spark.locality.wait", "0") \
.getOrCreate()
Кейс Б: Ритейл - real-time антифрод и LTV за 15 минут¶
Описание: крупный ритейлер, 15-минутные окна расчёта LTV на основе кликстрима. Объём: 50 TB активных данных. Требования: latency расчёта < 10 минут, высокая частота обновления, Spark Structured Streaming.
Анализ:
# Параметры:
active_data_tb = 50
batch_interval_min = 15 # каждые 15 минут новый расчёт
sla_latency_min = 10 # результат через 10 минут после события
# Паттерн доступа:
access_pattern = "hot" # данные читаются каждые 15 минут
streaming = True # Structured Streaming
shuffle_heavy = True # LTV = сложные join + aggregation
# Анализ требований:
# 1. Streaming каждые 15 минут → файлы пишутся часто → Small Files!
# 2. LTV = много join'ов → тяжёлый Shuffle
# 3. Latency SLA 10 минут → критично время планирования и чтения
# 4. 50 TB активных данных → можно хранить всё на HDFS
Рекомендация: HDFS с репликацией 3x для горячего слоя + ночная компакция.
- Streaming запись каждые 15 минут: HDFS избегает S3 API throttling ✅
- Тяжёлый Shuffle для LTV: NVMe HDFS превосходит S3 при intensive disk I/O ✅
- Data Locality: NODE_LOCAL для аналитических Stage'ов ✅
- Latency SLA: metadata операции в RAM NameNode (мс vs 50+ мс S3) ✅
Конфигурация Spark для этого кейса:
spark = SparkSession.builder \
.appName("retail-ltvCalculation-streaming") \
.config("spark.sql.streaming.schemaInference", "false") \
# HDFS HA конфигурация
.config("spark.hadoop.fs.defaultFS", "hdfs://retail-cluster") \
# Для тяжёлого Shuffle: локальные NVMe диски
.config("spark.local.dir", "/nvme/spark-local/") \
# Для streaming: не ждём локальность (микробатчи короткие)
.config("spark.locality.wait", "0") \
# Short-Circuit для максимальной скорости чтения
.config("spark.hadoop.dfs.client.read.shortcircuit", "true") \
.config("spark.hadoop.dfs.domain.socket.path",
"/var/lib/hadoop-hdfs/dn_socket") \
# AQE для adaptive coalescing после shuffle
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.getOrCreate()
# Ночная компакция горячего слоя
def compact_hot_layer(spark, hdfs_path: str, date_str: str) -> None:
"""
Компактирует мелкие файлы от Streaming в оптимальные Parquet-файлы.
Запускается в 02:00 каждую ночь через Airflow.
"""
df = spark.read.parquet(f"{hdfs_path}/date={date_str}/")
count = df.count()
# Целевой размер: 256 MB на файл
# При 50 KB/строку: ~5000 строк/файл
target_partitions = max(1, int(count / 50_000))
df.repartition(target_partitions) \
.write.mode("overwrite") \
.parquet(f"{hdfs_path}_compact/date={date_str}/")
Кейс В: Банк - DWH под жёстким контуром безопасности¶
Описание: крупный банк, корпоративное DWH. No-Cloud политика (данные не могут уходить в облако). Ночной расчёт регуляторной отчётности на Spark SQL - должен завершиться за 4 часа. Объём: 200 TB активных данных + 2 PB архива.
Анализ:
# Параметры:
active_tb = 200 # активных данных
archive_tb = 2_000 # архивных данных
no_cloud = True # требование безопасности
nightly_sla_hours = 4 # ночной расчёт за 4 часа
batch_type = "complex" # сложные join'ы, агрегации
# Стратегия:
# - Активные 200 TB: HDFS с репликацией 3x
# → максимальная производительность Spark SQL
# → data locality, Short-Circuit, низкий metadata latency
# - Архив 2 PB: Ceph RGW с S3 API + EC RS-6-3
# → 50% экономии vs HDFS replication
# → S3-совместимый API, но физически on-prem
# → Iceberg таблицы для прозрачного доступа
# Гибридная архитектура:
# HDFS (горячий) → Ceph S3 (архив)
# Оба on-prem, оба внутри банковского периметра
Рекомендация: гибридная архитектура HDFS + Ceph RGW с Iceberg.
- No-Cloud: оба хранилища on-prem ✅
- Активные данные (200 TB) на HDFS: максимум производительности для 4-часового SLA ✅
- Архив (2 PB) на Ceph с EC RS-6-3: 50% экономии дисков vs HDFS replication ✅
- Iceberg: единый каталог над HDFS и Ceph, прозрачные запросы через оба хранилища ✅
# Конфигурация Spark для банковского кейса
spark = SparkSession.builder \
.appName("bank-regulatory-reporting") \
# HDFS для горячих данных
.config("spark.hadoop.fs.defaultFS", "hdfs://bank-cluster") \
# Ceph RGW для архивных данных (через S3A)
.config("spark.hadoop.fs.s3a.endpoint", "http://ceph-rgw.bank.internal:7480") \
.config("spark.hadoop.fs.s3a.path.style.access", "true") \
.config("spark.hadoop.fs.s3a.aws.credentials.provider",
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider") \
# Iceberg с Hive Metastore (каталог над обоими хранилищами)
.enableHiveSupport() \
.config("spark.sql.hive.metastore.uris", "thrift://metastore:9083") \
.config("spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
# Производительность: максимум для ночного batch
.config("spark.executor.instances", "50") \
.config("spark.executor.cores", "8") \
.config("spark.executor.memory", "24g") \
# HDFS short-circuit для горячих данных
.config("spark.hadoop.dfs.client.read.shortcircuit", "true") \
.config("spark.hadoop.dfs.domain.socket.path",
"/var/lib/hadoop-hdfs/dn_socket") \
# Node locality для HDFS stage'ов
.config("spark.locality.wait.node", "15s") \
# AQE для адаптивной оптимизации
.config("spark.sql.adaptive.enabled", "true") \
.getOrCreate()
# Запрос, который прозрачно читает из обоих хранилищ
def run_regulatory_report(spark, report_date: str) -> None:
"""
Регуляторный отчёт ЦБ РФ N47-P.
Горячие данные - из HDFS, архивные - из Ceph, всё через Iceberg.
"""
# Текущие позиции (HDFS - горячие данные)
positions = spark.table("iceberg.bank_dw.positions") \
.filter(f"report_date = '{report_date}'")
# Исторические транзакции (Ceph - архив)
transactions_12m = spark.table("iceberg.bank_archive.transactions") \
.filter(f"transaction_date >= date_sub('{report_date}', 365)")
# Справочники рисков (HDFS)
risk_weights = spark.table("hive.bank_dw.risk_weights")
# Расчёт регуляторного норматива Н1
result = positions \
.join(transactions_12m, "account_id", "left") \
.join(risk_weights, "asset_class", "inner") \
.groupBy("legal_entity", "currency") \
.agg({
"risk_weighted_amount": "sum",
"capital_requirement": "sum",
})
result.write \
.mode("overwrite") \
.saveAsTable(f"iceberg.bank_reports.n1_report_{report_date.replace('-', '_')}")
Итоги: фреймворк принятия решений¶
После изучения всех аспектов сравнения - latency, consistency, atomicity, TCO, Committers, Lakehouse - можно сформулировать единый фреймворк.
Главный вывод урока: вопрос «HDFS или S3» давно перестал иметь однозначный ответ. Современный ответ - «зависит от контекста», и этот контекст определяется требованиями безопасности, характером нагрузки, экономической моделью и зрелостью инфраструктурной команды.
Lakehouse-форматы (Iceberg, Delta Lake) сделали S3 полноценной платформой для серьёзной аналитики, устранив исторические проблемы с LIST latency и атомарностью. При этом HDFS остаётся непревзойдённым для специфических сценариев: Heavy Shuffle, итеративный ML, strict data sovereignty, максимальная Spark-производительность при плотной нагрузке 24/7.
Архитектор данных должен понимать оба варианта достаточно глубоко, чтобы обосновать выбор - технически и экономически - перед бизнесом и командой.