Small Files в HDFS: нагрузка на NameNode heap и метаданные

Полный разбор проблемы Small Files в HDFS: физика метаданных NameNode Heap, математика катастрофы, GC паузы и ложные фейловеры, паралич планировщика Spark, AQE и maxPartitionBytes, стратегии записи, Daily Compaction и Hadoop Archives, диагностика через JMX и Spark UI.

storage platform

1. Физическая природа проблемы Small Files в архитектуре HDFS

Проблема мелких файлов - одна из самых распространённых причин деградации production HDFS-кластеров. Понять её природу невозможно без понимания того, как NameNode хранит метаданные файловой системы.

Анатомия метаданных NameNode: что хранится в Heap

NameNode - это не сервер хранения данных. Он не знает содержимого файлов и не передаёт их клиентам. Единственная функция NameNode - хранить пространство имён файловой системы: дерево директорий, список файлов в каждой директории, метаданные каждого файла и карту расположения каждого блока на DataNode.

Всё это хранится в оперативной памяти Java Heap NameNode в виде объектов:

  • INode - объект для каждой директории или файла. Содержит имя, владельца, группу, права доступа (rwxrwxrwx), время изменения, время доступа, флаг репликации/EC политики. Размер: ~150–200 байт.
  • BlockInfo - объект для каждого блока файла. Содержит идентификатор блока (Block ID), Generation Stamp, список трёх DataNode с репликами. Размер: ~150 байт.
  • DatanodeStorageInfo - информация о конкретном хранилище DataNode для каждой реплики блока. Размер: ~50 байт.

Суммарно на каждый файл в HDFS тратится: 1 INode + (размер_файла / block_size) × BlockInfo + реплики × DatanodeStorageInfo.

Схема показывает структуру объектов в Java Heap NameNode. Каждый файл - это INode + набор BlockInfo. При малом размере файлов количество INode-объектов растёт, а сами данные ничтожно малы.

Математика катастрофы: один файл vs миллион файлов

Рассмотрим два сценария с одинаковым объёмом данных - 1 TB:

Сценарий A: 1000 файлов по 1 GB (размер блока 256 MB)

Файлов: 1000
Блоков на файл: 1 GB / 256 MB = 4
Итого блоков: 1000 × 4 = 4000

Heap NameNode:
- INode для 1000 файлов: 1000 × 200 = 200 KB
- BlockInfo для 4000 блоков: 4000 × 150 = 600 KB
Итого: ~800 KB

Сценарий B: 10 000 000 файлов по 100 KB (мелкие файлы)

Файлов: 10 000 000
Блоков на файл: 1 (100 KB < 128 MB - один неполный блок)
Итого блоков: 10 000 000

Heap NameNode:
- INode для 10M файлов: 10M × 200 = 2 GB
- BlockInfo для 10M блоков: 10M × 150 = 1.5 GB
Итого: ~3.5 GB

При репликации 3x (DatanodeStorageInfo × 3):
- Дополнительно: 10M × 3 × 50 = 1.5 GB
ИТОГО: ~5 GB только на метаданные

Оба сценария хранят одинаковый объём данных (1 TB), но потребление Heap NameNode отличается в 6000 раз. При реальных масштабах (100 миллионов мелких файлов) NameNode требует 50–100 GB Heap только под метаданные.

Понятие Data-to-Metadata Ratio

Data-to-Metadata Ratio - отношение объёма полезных данных к объёму метаданных в NameNode. Это главный индикатор здоровья on-premise Data Lake:

# Мониторинг Data-to-Metadata Ratio через HDFS API
import subprocess
import json


def get_dfs_health_metrics(namenode_host: str, port: int = 9870) -> dict:
    """
    Получает ключевые метрики HDFS для оценки Data-to-Metadata Ratio.
    Чем выше отношение bytes_total / num_files_total - тем лучше.

    Рекомендуемое значение: > 50 MB на файл (средний размер файла).
    Критическое значение: < 1 MB на файл (кластер болен small files).
    """
    import urllib.request

    url = (f"http://{namenode_host}:{port}/jmx"
           f"?qry=Hadoop:service=NameNode,name=FSNamesystemState")

    with urllib.request.urlopen(url, timeout=10) as resp:
        data = json.loads(resp.read())

    for bean in data.get("beans", []):
        if "FSNamesystemState" not in bean.get("name", ""):
            continue

        num_files = bean.get("FilesTotal", 1)
        num_blocks = bean.get("BlocksTotal", 0)
        total_bytes = bean.get("CapacityUsed", 0)
        heap_used = bean.get("JvmMetrics.MemHeapUsedM", 0)

        avg_file_size_mb = (total_bytes / num_files / 1024 / 1024) if num_files > 0 else 0

        return {
            "files_total": num_files,
            "blocks_total": num_blocks,
            "capacity_used_tb": total_bytes / 1024**4,
            "avg_file_size_mb": avg_file_size_mb,
            "blocks_per_file": num_blocks / num_files if num_files > 0 else 0,
            "health": "CRITICAL" if avg_file_size_mb < 1
                      else "WARNING" if avg_file_size_mb < 32
                      else "OK",
        }


# Вывод метрик
metrics = get_dfs_health_metrics("namenode.example.com")
print(f"Файлов в HDFS: {metrics['files_total']:,}")
print(f"Средний размер файла: {metrics['avg_file_size_mb']:.1f} MB")
print(f"Здоровье кластера: {metrics['health']}")

# Типичные пороги:
# avg_file_size_mb > 128 MB  → ✅ Отличное состояние
# avg_file_size_mb 32-128 MB → ⚠️ Допустимо, но следует улучшить
# avg_file_size_mb 1-32 MB   → ⚠️ Предупреждение, требует внимания
# avg_file_size_mb < 1 MB    → ❌ Критично, срочная компакция!

2. Удар по мастер-ноде: нагрузка на NameNode Heap и GC Pauses

Когда Data-to-Metadata Ratio деградирует из-за накопления миллионов мелких файлов, NameNode начинает испытывать проблемы, которые каскадно распространяются на весь кластер.

Взрывное раздувание Heap: от гигабайт к сотням гигабайт

В небольшом кластере NameNode комфортно работает с 8–16 GB Heap. По мере роста Data Lake NameNode Heap конфигурируют на 64, 128 или 256 GB. Но неконтролируемое создание мелких файлов может заполнить даже 256 GB Heap.

Типичные источники «взрывного» роста числа файлов:

  • Kafka → HDFS стриминг: каждые 30 секунд Structured Streaming создаёт новый Parquet файл на каждый топик. При 50 топиках × 24 часа × 2 файла/минуту = 144 000 новых файлов в день.
  • partitionBy по high-cardinality колонке: .write.partitionBy("user_id") при 1 000 000 уникальных user_id создаёт 1 000 000 директорий и файлов за один Spark Job.
  • Многоуровневое партиционирование: .write.partitionBy("year", "month", "day", "hour", "event_type") при 8 типах событий × 24 часа = 192 партиции в день. За год - 70 000 директорий.

Кошмар Garbage Collection: Stop-The-World паузы и ложные фейловеры

NameNode - это Java-процесс. Java Garbage Collector периодически освобождает неиспользуемую память. При Heap размером 128–256 GB GC становится серьёзной проблемой.

Диаграмма показывает опасный каскад: Full GC → NameNode не отвечает → ZKFC инициирует failover → пока failover идёт, NameNode выходит из GC и «воскресает». Кластер попадает в непредсказуемое состояние.

Это не гипотетическая ситуация. В production HDFS-кластерах с сотнями миллионов файлов GC паузы длиной 30–120 секунд - реальная проблема, которая вызывает ложные failover'ы и приводит к паникующим звонкам в ночное время.

Конфигурация JVM для уменьшения GC давления:

# /etc/hadoop/conf/hadoop-env.sh - настройка JVM NameNode

# Для кластеров с < 50 миллионами файлов
export HADOOP_NAMENODE_OPTS="-Xms32g -Xmx32g \
  -XX:+UseG1GC \
  -XX:G1HeapRegionSize=32m \
  -XX:MaxGCPauseMillis=200 \
  -XX:+UnlockExperimentalVMOptions \
  -XX:G1NewSizePercent=10 \
  -XX:G1MaxNewSizePercent=25 \
  -XX:G1MixedGCLiveThresholdPercent=85 \
  -XX:+ParallelRefProcEnabled \
  -Xloggc:/var/log/hadoop/namenode-gc.log \
  -XX:+PrintGCDetails \
  -XX:+PrintGCDateStamps"

# Для кластеров с 100+ миллионами файлов (большой Heap)
# ZK failover timeout нужно ОБЯЗАТЕЛЬНО увеличить под GC паузы!
export HADOOP_NAMENODE_OPTS="-Xms128g -Xmx128g \
  -XX:+UseG1GC \
  -XX:G1HeapRegionSize=64m \
  -XX:MaxGCPauseMillis=1000 \
  -XX:SurvivorRatio=2 ..."

# В hdfs-site.xml увеличьте таймаут ZKFC под GC паузы:
# <property>
#   <name>ha.zookeeper.session-timeout.ms</name>
#   <value>60000</value>  ← 60 сек вместо дефолтных 10 сек
# </property>

Деградация RPC-очереди NameNode

Каждый запрос Spark Driver к NameNode (listStatus, getBlockLocations, create, rename) попадает в RPC Queue NameNode. При нормальной работе очередь обрабатывается быстро.

При миллионах мелких файлов ситуация меняется:

  • listStatus("/data/events/2024/") на директории с 500 000 файлами занимает секунды
  • Сотни параллельных Spark-приложений отправляют тысячи RPC-запросов одновременно
  • RPC Queue переполняется, новые запросы получают org.apache.hadoop.ipc.RetriableException
  • Spark Driver ждёт, Executor'ы простаивают, задания зависают
# Мониторинг RPC очереди NameNode через JMX
curl -s "http://namenode:9870/jmx?qry=Hadoop:service=NameNode,name=RpcActivityForPort8020" \
  | python3 -c "
import sys, json
data = json.load(sys.stdin)
for bean in data['beans']:
    if 'RpcActivity' not in bean.get('name', ''): continue
    print(f'RPC Queue Length:       {bean.get(\"RpcQueueLength\", 0)}')
    print(f'RPC Processing Time:    {bean.get(\"RpcProcessingTimeAvgTime\", 0):.1f} ms')
    print(f'RPC Queue Time:         {bean.get(\"RpcQueueTimeAvgTime\", 0):.1f} ms')
    print(f'Total RPC Calls:        {bean.get(\"RpcTotalCalls\", 0):,}')
    print()
    # Критичные пороги:
    rpc_queue = bean.get('RpcQueueLength', 0)
    if rpc_queue > 1000:
        print('КРИТИЧНО: RPC Queue переполнен! Вероятная причина: small files')
    elif rpc_queue > 100:
        print('ПРЕДУПРЕЖДЕНИЕ: RPC Queue нагружен')
"

3. Как Spark реагирует на мелкие файлы: паралич планировщика

Small Files в HDFS - это проблема не только NameNode. Spark-планировщик тоже страдает от миллионов мелких файлов, причём по своим специфическим причинам.

Антипаттерн «1 файл = 1 таска»

По умолчанию Spark при чтении файлов из HDFS создаёт одну Task на каждый InputSplit. Размер InputSplit по умолчанию = max(dfs.blocksize, spark.sql.files.maxPartitionBytes) = 128 MB.

Если файл меньше 128 MB, Spark создаёт один InputSplit для всего файла. Если файлов миллион - Tasks миллион.

Схема демонстрирует ключевое различие: при нормальных файлах количество Tasks пропорционально объёму данных. При мелких файлах количество Tasks определяется числом файлов, а не объёмом данных. Миллион задач для 1 GB данных - это абсурд.

Перегрузка Spark Driver: OOM при планировании

Прежде чем запустить хоть одну Task, Spark Driver должен:

  1. Опросить NameNode - получить блочные адреса для всех файлов
  2. Создать объект TaskDescription для каждой Task
  3. Сериализовать все TaskDescription для отправки Executor'ам
  4. Отслеживать статус каждой Task в памяти

При 1 000 000 Tasks каждый TaskDescription занимает ~200–500 байт → только описания тасок требуют 200–500 MB памяти Driver'а. Добавьте список всех файлов, блоков, предпочтений локальности - и Driver легко исчерпывает выделенную ему память.

# Типичный лог при Small Files проблеме в Spark Driver
INFO  DAGScheduler: Submitting 847392 missing tasks from ResultStage 0
  (MapPartitionsRDD) (first 15 tasks are for partitions
  Array(0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14))

WARN  TaskSetManager: Stage 0 contains a task of very large size (1526 KiB).
  The maximum recommended task size is 1000 KiB.

# Через несколько секунд:
ERROR SparkContext: Error while posting SparkListenerApplicationEnd
java.lang.OutOfMemoryError: Java heap space
    at org.apache.spark.scheduler.TaskDescription.decode(TaskDescription.scala:87)

Первая строка уже должна вас насторожить: 847 000 tasks для одного Stage. Это не аналитика - это дисфункция системы.

Потери времени на координацию: Scheduler Delay

В Spark UI на вкладке Stages метрика Scheduler Delay показывает время от момента готовности Task до её фактического запуска. При нормальной работе - миллисекунды. При Small Files картина иная:

Stage 0 (read from HDFS):
  Total Tasks:        847,392
  Duration:           47 minutes
  Scheduler Delay:    46 min 52 sec (99.4% времени!)
  Compute Time:       16 sec (0.6% времени)

Input Metrics:
  Records Read:       1,234,567
  Bytes Read:         2.3 GB

47 минут на чтение 2.3 GB данных - из которых 46 минут 52 секунды потрачено на overhead координации тасок. Само чтение и обработка данных заняли 16 секунд. Это классическая картина Small Files проблемы в Spark UI.


4. Встроенные спасатели Spark: File Max Split Size и AQE

Spark содержит несколько механизмов защиты от Small Files, которые нужно понимать и правильно настраивать.

Механизм Coalescing Inputs: как Spark объединяет мелкие файлы

Spark при чтении выполняет file coalescing - объединение нескольких маленьких файлов в один InputSplit. Это происходит в FilePartition.getFilePartitions() и управляется параметром openCostInBytes.

Алгоритм coalescing работает так:

  1. Берём список всех файлов, отсортированных по размеру
  2. Начинаем заполнять текущий InputSplit
  3. Для каждого файла добавляем его «штраф» за открытие (openCostInBytes) к его реальному размеру
  4. Пока суммарный размер Split < maxPartitionBytes - добавляем следующий файл в тот же Split
  5. Когда Split заполнен - начинаем новый
# Эффект openCostInBytes на coalescing
# Файлы: 100 файлов по 1 MB каждый
# maxPartitionBytes = 128 MB, openCostInBytes = 4 MB (дефолт)

# БЕЗ openCostInBytes (openCostInBytes = 0):
# Один InputSplit = 128 файлов × 1 MB = 128 MB → 1 Task
# Итого Tasks для 100 файлов: 1 Task (все в один split)

# С openCostInBytes = 4 MB:
# Effective file size = 1 MB + 4 MB = 5 MB (штраф за открытие)
# Один InputSplit = 128 MB / 5 MB ≈ 25 файлов → 1 Task
# Итого Tasks для 100 файлов: 4 Tasks (4 splits × 25 файлов)

# Почему openCostInBytes > 0 полезно?
# Spark не знает реальный размер файла до открытия.
# openCostInBytes - это штраф, моделирующий overhead на открытие
# файла: seek, handshake с DataNode, чтение метаданных Parquet footer.
# Он предотвращает создание слишком «жирных» Splits из тысяч файлов
# (что тоже неоптимально - один Task читает 1000 файлов подряд).

Ключевые параметры конфигурации чтения

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("small-files-tuning") \

    # ── Ключевой параметр: целевой размер одной Task при чтении ──────
    # Дефолт: 128 MB (spark.sql.files.maxPartitionBytes)
    # Spark объединяет файлы пока суммарный размер < maxPartitionBytes.
    # Для Small Files: увеличьте до 256 MB или 512 MB чтобы
    # уменьшить количество Tasks.
    # НО: слишком большой Split может замедлить конкретные задачи
    # если данные требуют heavy processing.
    .config("spark.sql.files.maxPartitionBytes", str(256 * 1024 * 1024)) \

    # ── Штраф за открытие файла (влияет на coalescing агрессивность) ─
    # Дефолт: 4 MB. Увеличьте для более агрессивного объединения.
    # При openCostInBytes = 128 MB: Spark объединяет файлы очень
    # агрессивно, даже если каждый файл по 100 KB.
    # Рекомендация для Small Files: 8-16 MB.
    .config("spark.sql.files.openCostInBytes", str(8 * 1024 * 1024)) \

    # ── Минимальное число партиций для coalesced чтения ──────────────
    # Если файлов мало (< этого числа), Spark не объединяет.
    # Дефолт: по числу ядер кластера. Обычно оставляют дефолт.
    # .config("spark.sql.files.minPartitionNum", "200") \

    .getOrCreate()

# Демонстрация эффекта
import os

# Читаем директорию с 10000 мелких файлов
df = spark.read.parquet("hdfs://cluster/data/small_files/")
print(f"Число партиций без оптимизации (дефолт 128MB): {df.rdd.getNumPartitions()}")
# Например: 2000 партиций для 10000 файлов по 25 MB

# С оптимизацией (256 MB):
spark.conf.set("spark.sql.files.maxPartitionBytes", str(256 * 1024 * 1024))
df_opt = spark.read.parquet("hdfs://cluster/data/small_files/")
print(f"Число партиций с 256 MB: {df_opt.rdd.getNumPartitions()}")
# Например: 1000 партиций - вдвое меньше Tasks

AQE: Adaptive Query Execution и coalescePartitions

AQE (Adaptive Query Execution) - механизм Spark, который адаптирует план выполнения во время самого выполнения, опираясь на реальную статистику. Для Small Files особенно важна функция coalescePartitions.

spark = SparkSession.builder \
    # Включить AQE (по умолчанию true в Spark 3.x)
    .config("spark.sql.adaptive.enabled", "true") \

    # Включить автоматическое объединение мелких shuffle партиций
    .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \

    # Целевой размер финальной объединённой партиции.
    # Дефолт: 64 MB. Для тяжёлых аналитических запросов увеличьте до 256 MB.
    .config("spark.sql.adaptive.advisoryPartitionSizeInBytes",
            str(256 * 1024 * 1024)) \

    # Минимальный размер партиции после объединения.
    # Партиции меньше этого размера будут объединены с соседними.
    # Дефолт: 1 MB. При Small Files можно увеличить до 32 MB.
    .config("spark.sql.adaptive.coalescePartitions.minPartitionSize",
            str(32 * 1024 * 1024)) \

    # Минимальное число выходных партиций после coalesce.
    # Предотвращает схлопывание в одну партицию при малом объёме данных.
    # Дефолт: 1.
    .config("spark.sql.adaptive.coalescePartitions.initialPartitionNum",
            "200") \

    .getOrCreate()

5. Проактивные стратегии борьбы на этапе записи из Spark

Лучшее лечение Small Files - профилактика. Правильные практики при записи данных позволяют избежать проблемы полностью.

Табу на high-cardinality partitionBy

partitionBy по колонкам с высокой кардинальностью - главный источник Small Files в production кластерах.

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder.getOrCreate()

df = spark.range(10_000_000).select(
    F.col("id").alias("user_id"),         # 10M уникальных значений!
    (F.rand() * 1000).cast("long").alias("product_id"),
    (F.rand() * 100).alias("amount"),
    F.current_timestamp().alias("event_time"),
)

# ❌ АНТИПАТТЕРН: partitionBy по user_id
# Создаёт 10 000 000 директорий и файлов!
# df.write.partitionBy("user_id").parquet("hdfs://cluster/bad/")

# ✅ ПРАВИЛЬНО: partitionBy по низкокардинальным колонкам
# Только дата (365 партиций в год) - разумно
df_with_date = df.withColumn(
    "event_date",
    F.date_format("event_time", "yyyy-MM-dd")
)
df_with_date.write \
    .partitionBy("event_date") \
    .parquet("hdfs://cluster/good/")
# Результат: 1 директория на дату, файлы в ней по 128-256 MB

# ✅ ДОПУСТИМО: двухуровневое партиционирование если кардинальность разумна
# date (365) × event_type (5) = 1825 партиций в год
df_with_type = df.withColumn(
    "event_date", F.date_format("event_time", "yyyy-MM-dd")
).withColumn(
    "event_type",
    F.array(
        F.lit("purchase"), F.lit("view"), F.lit("click"), F.lit("share"), F.lit("return")
    ).getItem((F.rand() * 5).cast("int"))
)

df_with_type.write \
    .partitionBy("event_date", "event_type") \
    .parquet("hdfs://cluster/acceptable/")

Рецепт правильного числа выходных файлов

Ключевой вопрос при записи: сколько файлов должно быть в результате? Оптимальный размер файла для HDFS - 256–512 MB (один блок или несколько полных блоков).

def calculate_optimal_partitions(
    df,
    target_file_size_mb: int = 256,
    sample_fraction: float = 0.01,
) -> int:
    """
    Вычисляет оптимальное число партиций для записи в HDFS.

    Алгоритм:
    1. Берём выборку из DataFrame (1%)
    2. Измеряем размер выборки в памяти
    3. Экстраполируем на весь DataFrame
    4. Делим на целевой размер файла
    5. Добавляем 20% запас для Parquet/Snappy сжатия

    target_file_size_mb: целевой размер файла в MB.
    256 MB = один HDFS блок (если blocksize=256MB).
    sample_fraction: доля данных для выборки.
    """
    sample = df.sample(fraction=sample_fraction, seed=42)

    # Материализуем выборку и измеряем размер
    # Через RDD.map(len).sum() измеряем сериализованный размер
    serialized_size = sample.rdd.map(lambda row: len(str(row))).sum()
    estimated_total_bytes = serialized_size / sample_fraction

    # Parquet со Snappy сжатием обычно в 5-10 раз меньше text
    # Для columnar Parquet используем коэффициент ~4
    parquet_estimated_bytes = estimated_total_bytes / 4

    target_bytes = target_file_size_mb * 1024 * 1024

    optimal_partitions = max(1, int(parquet_estimated_bytes / target_bytes * 1.2))

    print(f"Оценочный размер датасета: {parquet_estimated_bytes / 1024**3:.1f} GB (Parquet)")
    print(f"Целевой размер файла: {target_file_size_mb} MB")
    print(f"Рекомендуемое число партиций: {optimal_partitions}")

    return optimal_partitions


# Пример использования
df = spark.read.parquet("hdfs://cluster/input/large_dataset/")
n_partitions = calculate_optimal_partitions(df, target_file_size_mb=256)

df.repartition(n_partitions).write \
    .mode("overwrite") \
    .parquet("hdfs://cluster/output/optimized/")

coalesce vs repartition: когда что использовать

# coalesce(N): уменьшает число партиций БЕЗ full shuffle
# Работает путём объединения соседних партиций на тех же Executor'ах
# Не перемешивает данные, не изменяет физическое размещение блоков
# ОГРАНИЧЕНИЕ: работает только для УМЕНЬШЕНИЯ числа партиций

# ✅ ИСПОЛЬЗУЙТЕ coalesce когда:
# 1. Нужно уменьшить число партиций перед записью
# 2. Данные уже на правильных Executor'ах (после NODE_LOCAL чтения)
# 3. Не нужна равномерность распределения по файлам
df \
    .coalesce(50) \
    .write.parquet("hdfs://cluster/output/")
# Экономия: нет сетевого трафика Shuffle, быстрее


# repartition(N): перераспределяет данные через full Shuffle
# Создаёт N равномерных партиций с перемешиванием данных
# ДОРОЖЕ: требует записи и чтения Shuffle-файлов на диск

# ✅ ИСПОЛЬЗУЙТЕ repartition когда:
# 1. Нужна равномерность распределения данных по файлам
# 2. Перед partitionBy на high-cardinality колонке
# 3. Данные skewed и coalesce создаст неравномерные файлы
df \
    .repartition(50, F.col("event_date")) \  # repartition by column = hash-based
    .write.partitionBy("event_date").parquet("hdfs://cluster/output/")
# Каждая event_date-партиция получит ровно одинаковый файл


# maxRecordsPerFile: ограничение числа строк в одном файле
# Полезно когда размер строки трудно предсказать
# Не влияет на число партиций - только ограничивает размер файлов внутри партиции

spark.conf.set("spark.sql.maxRecordsPerFile", "1000000")  # 1M строк на файл

df.write \
    .partitionBy("event_date") \
    .parquet("hdfs://cluster/output/")
# Если в партиции date=2024-01-15 есть 5M строк:
# Без maxRecordsPerFile: 1 файл × 5M строк (возможно очень большой)
# С maxRecordsPerFile=1M: 5 файлов × 1M строк

Динамическое разбиение при partitionBy

Когда используется partitionBy в Spark, каждый Executor записывает данные в несколько партиций одновременно. При large shuffle это может создавать огромное число одновременно открытых файлов:

# Проблема: 200 Executor'ов × 365 date-партиций = 73000 открытых файлов
# Каждый файл минимум ~1 MB в памяти Executor'а → OOM

# Решение 1: предварительная сортировка перед partitionBy
# Каждый Executor записывает в одну партицию за раз
df \
    .sort("event_date") \
    .write \
    .partitionBy("event_date") \
    .parquet("hdfs://cluster/output/sorted/")
# Плюс: меньше открытых файлов одновременно
# Минус: дорогая сортировка (global sort = heavy Shuffle)

# Решение 2: repartition по партиционной колонке
# Каждый Executor отвечает за одну date-партицию
df \
    .repartition("event_date") \
    .write \
    .partitionBy("event_date") \
    .parquet("hdfs://cluster/output/repartitioned/")
# Плюс: один файл на date-партицию, нет смешивания
# Минус: полный Shuffle (как repartition)

# Настройка числа открытых файлов (Spark 3.x+)
spark.conf.set("spark.sql.sources.maxConcurrentWrites", "1")

6. Реактивные стратегии: Compaction Pipelines и Hadoop Archives

Проактивные меры помогают не создавать новые Small Files, но что делать с уже накопившимися миллионами мелких файлов?

Паттерн Daily Compaction: «Ночная приборка»

Daily Compaction - это фоновый процесс, который периодически (обычно ночью, когда кластер менее загружен) читает накопившиеся мелкие файлы за день и переписывает их в оптимальные крупные файлы.

# dags/daily_compaction_dag.py
# Airflow DAG для ночной компакции Small Files в HDFS

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
import logging

logger = logging.getLogger(__name__)


def compact_hdfs_partition(
    hdfs_path: str,
    date_str: str,
    target_file_size_mb: int = 256,
    output_format: str = "parquet",
) -> None:
    """
    Компактирует все мелкие файлы в указанной дате-партиции.

    Алгоритм «атомарной замены»:
    1. Читаем все файлы из hdfs_path/date=date_str/
    2. Записываем в hdfs_path_tmp/date=date_str/ (временная директория)
    3. Верифицируем: число строк в tmp == число строк в оригинале
    4. Переименовываем tmp → production (атомарно)
    5. Удаляем _SUCCESS файл из старой директории

    Критично: используем атомарный swap директорий чтобы
    ни в какой момент не было состояния "данные недоступны".
    """
    from pyspark.sql import SparkSession

    spark = SparkSession.builder \
        .appName(f"compaction-{date_str}") \
        .config("spark.sql.files.maxPartitionBytes",
                str(target_file_size_mb * 1024 * 1024)) \
        .getOrCreate()

    source_path = f"{hdfs_path}/date={date_str}"
    tmp_path = f"{hdfs_path}_compaction_tmp/date={date_str}"
    backup_path = f"{hdfs_path}_backup/date={date_str}"

    logger.info(f"Начинаем компакцию: {source_path}")

    # Шаг 1: читаем исходные данные
    df = spark.read.parquet(source_path)
    original_count = df.count()
    original_files = _count_files(source_path)

    logger.info(f"Исходно: {original_files} файлов, {original_count:,} строк")

    # Шаг 2: вычисляем оптимальное число выходных файлов
    optimal_partitions = calculate_optimal_partitions(df, target_file_size_mb)
    logger.info(f"Целевое число файлов: {optimal_partitions}")

    # Шаг 3: записываем в tmp директорию
    df.coalesce(optimal_partitions).write \
        .mode("overwrite") \
        .format(output_format) \
        .save(tmp_path)

    # Шаг 4: верификация
    df_tmp = spark.read.format(output_format).load(tmp_path)
    tmp_count = df_tmp.count()

    if tmp_count != original_count:
        raise ValueError(
            f"Верификация провалена! "
            f"Исходных строк: {original_count:,}, "
            f"После компакции: {tmp_count:,}. "
            f"Откатываем изменения."
        )

    logger.info(f"Верификация прошла: {tmp_count:,} строк")

    # Шаг 5: атомарный swap (backup → production)
    _hdfs_rename(source_path, backup_path)
    _hdfs_rename(tmp_path, source_path)
    _hdfs_delete(backup_path)

    result_files = _count_files(source_path)
    logger.info(
        f"Компакция завершена: {original_files}{result_files} файлов, "
        f"экономия: {(1 - result_files/original_files)*100:.0f}% файлов"
    )

    spark.stop()


def _count_files(hdfs_path: str) -> int:
    """Подсчёт файлов через hdfs dfs -count."""
    import subprocess
    result = subprocess.run(
        ["hdfs", "dfs", "-count", hdfs_path],
        capture_output=True, text=True
    )
    # Формат: DIR_COUNT  FILE_COUNT  CONTENT_SIZE  PATH
    parts = result.stdout.split()
    return int(parts[1]) if len(parts) >= 2 else 0


def _hdfs_rename(src: str, dst: str) -> None:
    import subprocess
    subprocess.run(["hdfs", "dfs", "-mv", src, dst], check=True)


def _hdfs_delete(path: str) -> None:
    import subprocess
    subprocess.run(["hdfs", "dfs", "-rm", "-r", "-skipTrash", path], check=True)


def calculate_optimal_partitions(df, target_file_size_mb: int) -> int:
    """Упрощённая версия из раздела 5."""
    try:
        sample_size = df.limit(10000).count()
        estimated_bytes = df.count() * 200  # ~200 байт/строка в Parquet
        target_bytes = target_file_size_mb * 1024 * 1024
        return max(1, int(estimated_bytes / target_bytes))
    except Exception:
        return 10  # дефолт при ошибке оценки


with DAG(
    dag_id="daily_hdfs_compaction",
    default_args={
        "owner": "data-platform",
        "retries": 1,
        "retry_delay": timedelta(minutes=30),
    },
    schedule_interval="0 3 * * *",  # каждый день в 3:00
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=["hdfs", "compaction", "small-files"],
) as dag:

    from airflow.operators.python import PythonOperator
    from datetime import date, timedelta

    compact_yesterday = PythonOperator(
        task_id="compact_bronze_yesterday",
        python_callable=compact_hdfs_partition,
        op_kwargs={
            "hdfs_path": "hdfs://cluster/data/bronze/events",
            "date_str": "{{ ds }}",    # Airflow date macro: вчерашняя дата
            "target_file_size_mb": 256,
        },
        doc_md="""
        Компактирует Bronze-слой за вчерашний день.
        Уменьшает число файлов, освобождает Heap NameNode,
        ускоряет последующие Spark запросы.
        """,
    )

Hadoop Archives (HAR): архивирование миллионов старых файлов

Для очень старых данных, которые нужно хранить для compliance но почти никогда не читать, существует механизм Hadoop Archives (HAR). Он упаковывает тысячи файлов в один HAR-архив, резко снижая число объектов в NameNode.

# Создание HAR архива из директории с мелкими файлами
# hadoop archive: упаковывает исходную директорию в один .har файл
# -archiveName: имя архива (должно заканчиваться на .har)
# -p: parent директория (относительно неё будут пути внутри архива)
# /hdfs/source/path: что архивировать
# /hdfs/dest/path: куда положить архив

hadoop archive \
  -archiveName events_2020.har \
  -p /data/bronze/events \
  date=2020-01-01 date=2020-01-02 date=2020-01-03 \
  /data/archive/

# Результат: создан /data/archive/events_2020.har/
# Внутри архива: виртуальная структура директорий + индексный файл
# На HDFS: только 3-4 больших файла вместо тысяч мелких!

# Проверить содержимое HAR архива
hdfs dfs -ls har:///data/archive/events_2020.har/date=2020-01-01/

# Чтение через Spark прозрачно через har:// протокол
df = spark.read.parquet("har:///data/archive/events_2020.har/date=2020-01-01/")
df.show()

# Удаляем исходные файлы (они теперь в архиве)
hdfs dfs -rm -r /data/bronze/events/date=2020-01-01
hdfs dfs -rm -r /data/bronze/events/date=2020-01-02
hdfs dfs -rm -r /data/bronze/events/date=2020-01-03

Ограничения HAR:

  • HAR архивы не поддерживают изменение: нельзя добавить файл в существующий архив
  • Производительность чтения хуже, чем прямых Parquet файлов: HAR читается через индексный файл с дополнительным overhead'ом
  • HAR нельзя использовать как источник для stриминга или частого intra-day обращения
  • Рекомендуется для данных старше 3-5 лет, которые читаются не чаще раза в квартал

Hive Metastore и Iceberg Catalog как абстракция над файлами

Одна из корневых причин, почему Small Files так болезненно влияют на Spark - это необходимость обходить физические файлы через listStatus() NameNode при каждом чтении. Если Spark знает где файлы через каталог метаданных, можно избежать дорогостоящего listing'а.

# С Hive Metastore Spark читает список файлов из HMS, не из HDFS listing
# Это критически снижает нагрузку на NameNode RPC при small files

spark = SparkSession.builder \
    .enableHiveSupport() \
    .config("spark.sql.hive.metastore.uris", "thrift://metastore:9083") \
    .getOrCreate()

# Вместо прямого пути (вызывает listStatus на весь путь):
# df = spark.read.parquet("hdfs://cluster/data/events/")

# Читаем через Hive таблицу (HMS знает точный список файлов):
df = spark.table("analytics.events")
# Spark делает запрос в HMS: "дай мне файлы для table=analytics.events,
# date=2024-01-15" и получает точный список без NameNode listing!

# Apache Iceberg даёт ещё большую независимость от NameNode:
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") \
    .getOrCreate()

# Iceberg хранит точный manifest list файлов в metadata/
# Никакого listStatus() NameNode при чтении!
df = spark.table("iceberg.analytics.events")
df.filter("event_date = '2024-01-15'").show()
# Iceberg знает точно какие Parquet файлы содержат date=2024-01-15
# и обращается напрямую - без сканирования директорий

7. Практика: аудит NameNode и Spark UI, лабораторный бенчмарк

Диагностика: как обнаружить Small Files проблему

Полный диагностический скрипт

# small_files_audit.py
# Комплексная диагностика Small Files проблемы в HDFS + Spark

import subprocess
import json
import urllib.request
from dataclasses import dataclass
from typing import Optional


@dataclass
class HDFSHealthReport:
    files_total: int
    blocks_total: int
    capacity_used_tb: float
    avg_file_size_mb: float
    heap_used_pct: float
    gc_count_major: int
    rpc_queue_length: int
    health_score: str


def get_namenode_metrics(namenode_host: str, port: int = 9870) -> dict:
    """Собирает все нужные метрики с NameNode через JMX HTTP API."""
    metrics = {}

    # FSNamesystemState: статистика файлов и блоков
    url = f"http://{namenode_host}:{port}/jmx?qry=Hadoop:service=NameNode,name=FSNamesystemState"
    try:
        with urllib.request.urlopen(url, timeout=10) as resp:
            data = json.loads(resp.read())
        for bean in data.get("beans", []):
            if "FSNamesystemState" not in bean.get("name", ""):
                continue
            metrics.update({
                "files_total": bean.get("FilesTotal", 0),
                "blocks_total": bean.get("BlocksTotal", 0),
                "capacity_used": bean.get("CapacityUsed", 0),
                "capacity_total": bean.get("CapacityTotal", 1),
            })
    except Exception as e:
        metrics["error_fsnamesystem"] = str(e)

    # JvmMetrics: использование Heap и GC статистика
    url2 = f"http://{namenode_host}:{port}/jmx?qry=Hadoop:service=NameNode,name=JvmMetrics"
    try:
        with urllib.request.urlopen(url2, timeout=10) as resp:
            data2 = json.loads(resp.read())
        for bean in data2.get("beans", []):
            if "JvmMetrics" not in bean.get("name", ""):
                continue
            heap_used = bean.get("MemHeapUsedM", 0)
            heap_max = bean.get("MemHeapMaxM", 1)
            metrics.update({
                "heap_used_mb": heap_used,
                "heap_max_mb": heap_max,
                "heap_used_pct": heap_used / heap_max * 100,
                "gc_count": bean.get("GcCount", 0),
                "gc_time_millis": bean.get("GcTimeMillis", 0),
            })
    except Exception as e:
        metrics["error_jvm"] = str(e)

    # RPC метрики
    url3 = (f"http://{namenode_host}:{port}/jmx"
            f"?qry=Hadoop:service=NameNode,name=RpcActivityForPort8020")
    try:
        with urllib.request.urlopen(url3, timeout=10) as resp:
            data3 = json.loads(resp.read())
        for bean in data3.get("beans", []):
            if "RpcActivity" not in bean.get("name", ""):
                continue
            metrics.update({
                "rpc_queue_length": bean.get("RpcQueueLength", 0),
                "rpc_processing_avg_ms": bean.get("RpcProcessingTimeAvgTime", 0),
                "rpc_queue_avg_ms": bean.get("RpcQueueTimeAvgTime", 0),
            })
    except Exception as e:
        metrics["error_rpc"] = str(e)

    return metrics


def generate_health_report(namenode_host: str) -> HDFSHealthReport:
    """Генерирует полный отчёт о здоровье HDFS с оценкой."""
    m = get_namenode_metrics(namenode_host)

    files_total = m.get("files_total", 0)
    capacity_used = m.get("capacity_used", 0)
    avg_file_size_mb = (capacity_used / files_total / 1024 / 1024
                        if files_total > 0 else 0)
    heap_used_pct = m.get("heap_used_pct", 0)
    rpc_queue = m.get("rpc_queue_length", 0)

    # Оценка здоровья
    issues = []
    if avg_file_size_mb < 1:
        issues.append("КРИТИЧНО: средний размер файла < 1 MB")
    elif avg_file_size_mb < 32:
        issues.append("ПРЕДУПРЕЖДЕНИЕ: средний размер файла < 32 MB")

    if files_total > 500_000_000:
        issues.append("КРИТИЧНО: > 500M файлов в NameNode")
    elif files_total > 100_000_000:
        issues.append("ПРЕДУПРЕЖДЕНИЕ: > 100M файлов в NameNode")

    if heap_used_pct > 85:
        issues.append("КРИТИЧНО: Heap NameNode > 85% заполнен")
    elif heap_used_pct > 70:
        issues.append("ПРЕДУПРЕЖДЕНИЕ: Heap NameNode > 70% заполнен")

    if rpc_queue > 500:
        issues.append("КРИТИЧНО: RPC Queue > 500 запросов")
    elif rpc_queue > 100:
        issues.append("ПРЕДУПРЕЖДЕНИЕ: RPC Queue > 100 запросов")

    health_score = "CRITICAL" if any("КРИТИЧНО" in i for i in issues) else \
                   "WARNING" if issues else "OK"

    report = HDFSHealthReport(
        files_total=files_total,
        blocks_total=m.get("blocks_total", 0),
        capacity_used_tb=capacity_used / 1024**4,
        avg_file_size_mb=avg_file_size_mb,
        heap_used_pct=heap_used_pct,
        gc_count_major=m.get("gc_count", 0),
        rpc_queue_length=rpc_queue,
        health_score=health_score,
    )

    # Вывод отчёта
    print("\n" + "=" * 60)
    print("ОТЧЁТ О ЗДОРОВЬЕ HDFS (Small Files аудит)")
    print("=" * 60)
    print(f"Файлов в NameNode:     {report.files_total:>15,}")
    print(f"Блоков в NameNode:     {report.blocks_total:>15,}")
    print(f"Использовано:          {report.capacity_used_tb:>14.1f} TB")
    print(f"Средний размер файла:  {report.avg_file_size_mb:>14.1f} MB")
    print(f"Heap NameNode:         {report.heap_used_pct:>14.1f}%")
    print(f"RPC Queue:             {report.rpc_queue_length:>15,}")
    print()
    print(f"Оценка здоровья: {health_score}")

    if issues:
        print("\nПроблемы:")
        for issue in issues:
            print(f"  ⚠️  {issue}")
    else:
        print("✅ Все метрики в норме")

    return report


# Пример запуска
report = generate_health_report("namenode.example.com")

Лабораторный бенчмарк: «раздробленная» vs «компактная» таблица

# lab_small_files_benchmark.py
# Демонстрирует влияние Small Files на производительность Spark

import time
from pyspark.sql import SparkSession, functions as F


def create_spark() -> SparkSession:
    return (
        SparkSession.builder
        .master("yarn")
        .appName("small-files-benchmark")
        .config("spark.executor.instances", "10")
        .config("spark.executor.cores", "4")
        .config("spark.executor.memory", "8g")
        .config("spark.sql.adaptive.enabled", "false")  # отключаем AQE для чистоты
        .getOrCreate()
    )


def create_fragmented_dataset(spark: SparkSession, output_path: str) -> None:
    """
    Создаёт датасет из 50000 мелких файлов.
    Имитирует типичный результат Streaming-записи без компакции.
    """
    print(f"Создаём 'раздробленный' датасет → {output_path}")

    for i in range(50000):
        # Каждая запись создаёт отдельный файл через coalesce(1)
        df = spark.range(100).select(
            F.col("id"),
            F.lit(f"event_{i % 10}").alias("event_type"),
            F.rand().alias("amount"),
            F.current_timestamp().alias("ts"),
        )
        df.coalesce(1).write \
            .mode("append") \
            .parquet(f"{output_path}/micro_batch={i:05d}/")

    print(f"Создано 50000 файлов, ~100 строк каждый")


def create_compact_dataset(spark: SparkSession,
                            fragmented_path: str, compact_path: str) -> None:
    """Компактирует раздробленный датасет в оптимальные файлы."""
    print(f"Компактируем → {compact_path}")

    spark.conf.set("spark.sql.adaptive.enabled", "true")

    df = spark.read.parquet(fragmented_path)
    total_rows = df.count()

    # repartition(50): ~100K строк на файл при 5M строк всего
    df.repartition(50).write \
        .mode("overwrite") \
        .parquet(compact_path)

    print(f"Компактировано: {total_rows:,} строк в 50 файлов")


def run_analytical_query(spark: SparkSession,
                          data_path: str, label: str) -> float:
    """
    Типичный аналитический запрос: группировка + агрегация.
    Возвращает время выполнения в секундах.
    """
    print(f"\nЗапрос к: {label}")

    start = time.time()
    df = spark.read.parquet(data_path)

    result = df \
        .groupBy("event_type") \
        .agg(
            F.count("*").alias("cnt"),
            F.sum("amount").alias("total"),
            F.avg("amount").alias("avg"),
        )

    count = result.count()
    elapsed = time.time() - start

    print(f"  Строк: {count}, время: {elapsed:.1f} сек")
    print(f"  Task count: {spark.sparkContext.statusTracker().getActiveStageIds()}")

    return elapsed


if __name__ == "__main__":
    spark = create_spark()

    FRAGMENTED = "hdfs://cluster/lab/fragmented"
    COMPACT = "hdfs://cluster/lab/compact"

    # Создаём тестовые данные
    create_fragmented_dataset(spark, FRAGMENTED)
    create_compact_dataset(spark, FRAGMENTED, COMPACT)

    # Бенчмарк
    t_fragmented = run_analytical_query(spark, FRAGMENTED, "Раздробленные файлы (50000)")
    t_compact = run_analytical_query(spark, COMPACT, "Компактные файлы (50)")

    print(f"\n{'='*50}")
    print("РЕЗУЛЬТАТЫ")
    print(f"{'='*50}")
    print(f"Раздробленные (50000 файлов): {t_fragmented:.1f} сек")
    print(f"Компактные    (50 файлов):    {t_compact:.1f} сек")
    print(f"Ускорение: {t_fragmented / t_compact:.1f}x")
    print()
    print("Проверьте в Spark UI:")
    print("  - Число Tasks в каждом Stage")
    print("  - Scheduler Delay (должен быть значительно выше для раздробленных)")
    print("  - Input Metrics → Files Count")

    spark.stop()

Итоги: правила работы с HDFS для Data Engineer

Small Files - это проблема, которая незаметно нарастает и становится катастрофой в самый неподходящий момент. Несколько практических правил помогут избежать этого:

Правило 1: Знайте Data-to-Metadata Ratio своего кластера. Средний размер файла < 32 MB - сигнал тревоги. < 1 MB - срочное вмешательство.

Правило 2: Никогда не используйте partitionBy по high-cardinality колонкам. user_id, transaction_id, timestamp как ключи партиционирования - прямая дорога к миллионам файлов.

Правило 3: Настраивайте число выходных файлов явно. Всегда завершайте пайплайн записи через coalesce(N) или repartition(N) с осознанным N. Никогда не оставляйте Spark создавать произвольное число файлов.

Правило 4: Развёртывайте Daily Compaction для стриминговых слоёв. Structured Streaming создаёт мелкие файлы по природе. Ночная компакция - стандартная практика для production Data Lake.

Правило 5: Включайте AQE. spark.sql.adaptive.coalescePartitions.enabled = true - это бесплатная защита от мелких Shuffle-партиций внутри Job'ов.

Правило 6: Используйте Hive Metastore или Iceberg. Каталог метаданных снижает нагрузку на NameNode при leaf-file discovery - Spark знает где искать, не сканируя директории.