HDFS vs S3: когда что выбирать - latency, consistency, стоимость

Архитектурное сравнение HDFS и S3-совместимых хранилищ для Data Engineer: философия Co-located vs Decoupled, битва latency и throughput, гарантии consistency и атомарность rename, влияние на Spark Committers и Lakehouse форматы, модель TCO, матрица выбора и три кейс-стади из enterprise-практики.

storage platform

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 (объект сразу виден)
  • PUT UPDATE существующего объекта → 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 - это:

  1. GET list всех объектов с префиксом tmp/output_tmp/
  2. PUT каждый объект под новым ключом data/output/...
  3. 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 проявляется в:

  1. Scale-to-Zero: вычислительный кластер существует только во время обработки. HDFS требует постоянно работающих DataNode серверов.
  2. Cold Storage tiering: S3 Glacier (3-5x дешевле S3 Standard) для архивных данных. HDFS не имеет встроенного tiering.
  3. Отсутствие 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.

Архитектор данных должен понимать оба варианта достаточно глубоко, чтобы обосновать выбор - технически и экономически - перед бизнесом и командой.