Small Files в HDFS: нагрузка на NameNode heap и метаданные
Полный разбор проблемы Small Files в HDFS: физика метаданных NameNode Heap, математика катастрофы, GC паузы и ложные фейловеры, паралич планировщика Spark, AQE и maxPartitionBytes, стратегии записи, Daily Compaction и Hadoop Archives, диагностика через JMX и Spark UI.
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 должен:
- Опросить NameNode - получить блочные адреса для всех файлов
- Создать объект
TaskDescriptionдля каждой Task - Сериализовать все TaskDescription для отправки Executor'ам
- Отслеживать статус каждой 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 работает так:
- Берём список всех файлов, отсортированных по размеру
- Начинаем заполнять текущий InputSplit
- Для каждого файла добавляем его «штраф» за открытие (
openCostInBytes) к его реальному размеру - Пока суммарный размер Split <
maxPartitionBytes- добавляем следующий файл в тот же Split - Когда 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 знает где искать, не сканируя директории.