Small Files Disease: причины, последствия и влияние на планировщик
Почему 1000 файлов по 1 КБ медленнее одного файла в 1 МБ, как Driver падает с OOM от метаданных, откуда берётся болезнь в streaming и CDC пайплайнах, и как лечить через Compaction, AQE и современные Lakehouse форматы.
Есть проблема, которая встречается в каждом production Spark-кластере, который существует достаточно долго. Её не видно сразу - она нарастает незаметно: месяц за месяцем пайплайны пишут мелкие файлы, и однажды запрос, который раньше занимал 30 секунд, начинает висеть 40 минут. Driver падает с OutOfMemoryError на стадии «Listing leaf files and directories». Администраторы HDFS приходят с графиком, где NameNode memory растёт вертикально. В S3 счёт за хранение вырос вдвое, хотя данных столько же.
Это Small Files Disease - болезнь мелких файлов. В этом уроке разберём физику того, почему размер файла критически важен, как именно Spark-планировщик деградирует от миллионов мелких файлов, как диагностировать болезнь по симптомам в Spark UI и как лечить её на разных уровнях - от конфигурации до архитектурных решений в Delta Lake и Iceberg.
1. Что такое Small Files Disease и откуда она берётся¶
Определение: какой файл считается «мелким»?¶
Понятие «мелкий» относительно и зависит от контекста хранилища. Общепринятые ориентиры:
- Оптимальный размер файла в Lakehouse: 128 МБ - 512 МБ на файл в Parquet/ORC
- Нижняя граница нормы: 32–64 МБ - ещё приемлемо, но уже не идеально
- «Мелкий» файл: менее 16–32 МБ - начинают проявляться симптомы
- «Микроскопический» файл: менее 1 МБ - тяжёлые симптомы на любом масштабе
- Катастрофа: файлы по несколько КБ - кластер падает задолго до обработки данных
Почему именно 128–512 МБ? Это связано с размером HDFS-блока (128 МБ по умолчанию), оптимальным размером Parquet Row Group (128 МБ) и балансом между параллелизмом (один файл = один Task минимум) и overhead на открытие файла.
# Проверить средний размер файлов в таблице
import os
def check_file_sizes(path):
"""Диагностика размера файлов через os.walk."""
sizes = []
for root, dirs, files in os.walk(path):
for f in files:
if f.endswith((".parquet", ".orc")):
full_path = os.path.join(root, f)
sizes.append(os.path.getsize(full_path))
if not sizes:
print("Файлы не найдены")
return
avg_mb = sum(sizes) / len(sizes) / (1024 * 1024)
min_mb = min(sizes) / (1024 * 1024)
max_mb = max(sizes) / (1024 * 1024)
print(f"Файлов: {len(sizes):,}")
print(f"Средний размер: {avg_mb:.2f} МБ")
print(f"Минимальный: {min_mb:.3f} МБ")
print(f"Максимальный: {max_mb:.2f} МБ")
if avg_mb < 16:
print("КРИТИЧНО: Small Files Disease в тяжёлой форме")
elif avg_mb < 64:
print("ВНИМАНИЕ: файлы меньше оптимального размера")
else:
print("Норма: размер файлов в допустимом диапазоне")
check_file_sizes("/tmp/events_by_date/")
Причина 1: Неконтролируемый параллелизм при записи¶
Это самая распространённая причина. Когда Spark записывает данные, каждый активный Task создаёт отдельный файл в каждой partition-папке. Если у вас 200 Tasks (по умолчанию spark.sql.shuffle.partitions = 200) и таблица содержит 1 ГБ данных - каждый Task пишет файл размером 5 МБ. Если данных всего 50 МБ - каждый Task пишет файл по 250 КБ.
from pyspark.sql import functions as F
# Демонстрация: 50 МБ данных, 200 Tasks → 200 файлов по ~250 КБ
small_df = spark.range(1, 100_001) \
.withColumn("value", F.rand()) \
.withColumn("dt", F.lit("2026-01-15"))
# Количество shuffle partitions по умолчанию: 200
print(spark.conf.get("spark.sql.shuffle.partitions")) # → 200
# Запись: Spark создаст ~200 файлов
small_df.repartition(200).write \
.mode("overwrite") \
.partitionBy("dt") \
.parquet("/tmp/small_files_demo/")
# Итог: 200 файлов по ~250 КБ в папке /tmp/small_files_demo/dt=2026-01-15/
# Хотя данных всего ~8 МБ!
Математика простая: num_files = num_tasks × num_partition_values. При 200 Tasks и 365 днях партиционирования - это 73 000 файлов за один запуск. При ежедневных запусках - миллионы файлов за год.
Причина 2: Злоупотребление партиционированием (высокая кардинальность)¶
Если partitionBy выбран по колонке с высокой кардинальностью - в каждой partition-папке будет очень мало строк, а значит файлы будут крошечными. Эту причину мы разобрали в предыдущем уроке: высокая кардинальность → много маленьких папок → маленькие файлы в каждой.
Причина 3: Structured Streaming без компакции¶
Spark Structured Streaming записывает данные micro-batch за micro-batch. Если интервал micro-batch - 10 секунд, и каждый batch пишет 2–5 файлов, то за сутки накапливается:
24 × 60 × 6 batches/мин × 3 файла = 25 920 файлов в день
За месяц - ~778 000 файлов. За год - ~9.5 миллионов. И это только для одной таблицы.
# Streaming пишет micro-batches - каждый trigger создаёт новые файлы
# query = events_stream \
# .writeStream \
# .trigger(processingTime="10 seconds") \
# .partitionBy("event_date") \
# .format("parquet") \
# .option("path", "s3a://data/events/") \
# .option("checkpointLocation", "s3a://checkpoints/events/") \
# .start()
# Через 1 час: 360 micro-batches × 3 файла = 1 080 файлов
# Через 1 день: 8 640 micro-batches × 3 файла = 25 920 файлов
# Без компакции это накапливается бесконечно
Причина 4: CDC (Change Data Capture) пайплайны¶
CDC пайплайны реплицируют изменения из operational БД в lakehouse. Каждое изменение (INSERT, UPDATE, DELETE) - маленькое событие. Если CDC пишет в lakehouse напрямую без буферизации, каждый commit может добавлять один-два файла с несколькими тысячами строк.
2. Влияние на хранилище: как мелкие файлы убивают storage tier¶
HDFS: смерть NameNode по памяти¶
В HDFS каждый файл и каждый блок - это объект в оперативной памяти NameNode. NameNode хранит полный namespace (дерево директорий) и mapping «файл → список блоков» целиком в heap. Это необходимо для быстрого доступа - но это же делает NameNode узким местом при работе с миллионами мелких файлов.
Каждый объект в namespace занимает ~150–300 байт памяти NameNode:
- inode (файл или директория): ~150–180 байт
- block descriptor (для файла ≤128 МБ): ~112 байт
- итого для одного мелкого файла: ~300 байт в RAM NameNode
Пример: 10 миллионов мелких файлов
→ 10 000 000 × 300 байт = 3 ГБ только для метаданных
→ NameNode с 16 ГБ heap может хранить ≈53 миллиона объектов
→ Один большой кластер с плохим партиционированием легко достигает предела
→ Далее: GC паузы, деградация, падение NameNode
При деградации NameNode страдает не только ваш пайплайн - Directory Listing замедляется для всех пользователей кластера. GC Stop-The-World паузы блокируют NameNode на секунды. В экстремальных случаях NameNode падает с OOM - и весь кластер становится недоступен.
S3 и объектные хранилища: финансовая и операционная катастрофа¶
В объектных хранилищах нет NameNode - метаданные распределены по серверам провайдера. Но проблема мелких файлов проявляется через стоимость API-запросов и Rate Limits.
Финансовая сторона в AWS S3:
LIST Objects: $0.005 за 1000 запросовGET Object(чтение файла): $0.0004 за 1000 запросовPUT Object(запись файла): $0.005 за 1000 запросов
Пример: читаем таблицу из 10 миллионов мелких файлов
Directory Listing:
10 000 000 файлов / 1000 per request = 10 000 LIST запросов
Стоимость: 10 × $0.05 = $0.05
Чтение файлов:
10 000 000 GET запросов
Стоимость: 10 000 × $0.0004 = $4.00
Итого за ОДИН запрос к таблице: ~$4.05
При 100 запросах в день: $405 в день = $12 150 в месяц
И это только стоимость API-запросов, не считая трафик и хранение!
При тех же 1 ТБ данных, но в 1000 файлах по 1 ГБ:
1000 GET запросов = $0.0004 - в 10 000 раз дешевле!
Дополнительная проблема - Rate Limits: AWS ограничивает GET/HEAD до 5500 запросов/с на «префикс» (папку). При 10 000 параллельных Tasks, каждый из которых делает несколько GET-запросов, легко достичь лимита и получить ошибки 503 SlowDown. Spark будет их ретраить, ещё больше замедляя работу.
3. Под капотом Spark: что происходит с планировщиком¶
Страдания Driver: Directory Listing и OOM¶
Перед тем как Spark запустит хоть один Task, Driver обязан:
- Собрать список всех файлов для чтения (Directory Listing)
- Определить, какие partition-папки нужны (Partition Pruning)
- Разбить файлы на InputSplits для Tasks
- Создать объект TaskDescription для каждого Task
- Сериализовать и отправить TaskDescriptions на Executor-ы
Шаги 1–3 выполняются полностью на Driver в single-thread режиме. Все пути к файлам хранятся в памяти Driver в виде Java-объектов. При миллионе файлов:
1 000 000 путей × ~200 байт (String объект в JVM) = 200 МБ только для путей
+ metadata объекты FileStatus: ещё ~500 байт на файл
= 700 МБ в heap Driver только для метаданных файлов
При 10 миллионах файлов: ~7 ГБ
Типичный heap Driver: 4-8 ГБ → OutOfMemoryError
Характерный симптом в логах:
INFO FileSourceScanExec: Listing leaf files and directories...
(зависает на этой строке на несколько минут)
ERROR SparkContext: Error initializing SparkContext
java.lang.OutOfMemoryError: GC overhead limit exceeded
Взрывное планирование Tasks: бюрократия дороже работы¶
Даже если Driver справился с листингом и не упал по OOM, следующая проблема - количество Tasks. Один файл ≥ один Task (мелкие файлы никогда не объединяются без специальных настроек). Каждый Task - это:
- Создание объекта задачи на Driver: 1–5 мс
- Сериализация TaskDescription: 2–10 мс
- Передача по сети на Executor: 5–20 мс
- Десериализация на Executor: 2–10 мс
- Инициализация ресурсов (buffers, readers): 5–20 мс
- Реальная работа (прочитать 100 КБ): 1–2 мс
- Отчёт о завершении обратно на Driver: 2–5 мс
Суммарный overhead на Task: 17–70 мс. Реальная работа: 1–2 мс. КПД: 2–10%.
Оба сценария обрабатывают одинаковый объём данных (1 ГБ), но скорость отличается в 24 раза. Это и есть суть Small Files Disease: данных мало, работы мало, но накладные расходы убивают производительность.
Разрушение Parquet-оптимизаций¶
Parquet - колоночный формат, оптимизированный для чтения больших блоков последовательно. Его эффективность зависит от размера Row Group (обычно 128 МБ). При мелких файлах эта оптимизация разрушается:
- Vectorized читатель не может работать эффективно: нет достаточно большого batch
- Predicate Pushdown деградирует: min/max статистики Row Group менее эффективны при 1000 строках вместо 1 миллиона
- Фиксированный overhead при открытии файла (header, footer, column index) составляет ~5–20 КБ. При файле 100 КБ это 5–20% overhead. При файле 128 МБ - 0.004–0.016%
- Compression ratio хуже: сжатие лучше работает на больших однородных блоках данных
4. Диагностика: как обнаружить Small Files Disease¶
Симптомы в Spark UI¶
Открываем Spark UI → Stages и смотрим на стадию чтения данных.
Сигнал тревоги 1: огромное число Tasks при небольшом Input Size. Если стадия чтения запустила 20 000 Tasks, а общий Input Size составляет 2 ГБ - каждый Task читает в среднем 100 КБ. Это явный признак мелких файлов.
Сигнал тревоги 2: Task Duration очень мал, но Scheduler Delay велик. В Timeline каждого Task есть цветная полоска:
- Синий (Scheduler Delay): время от создания Task до начала выполнения
- Жёлтый (Task Deserialization Time): время на десериализацию
- Зелёный (Executor Computing Time): реальная работа
- Красный (Result Serialization Time): сериализация результата
При Small Files Disease: зелёная полоска (реальная работа) крошечная - 1–5 мс. Синяя и жёлтая занимают 80–95% времени каждого Task.
Сигнал тревоги 3: долгая стадия «Listing leaf files». В описании Job/Stage появляется «Listing leaf files and directories for N paths» - и он висит минутами.
# Программная диагностика: смотреть на число partitions при чтении
from pyspark.sql import functions as F
df = spark.read.parquet("/path/to/table/")
num_partitions = df.rdd.getNumPartitions()
print(f"Число Spark-партиций (≈ число файлов): {num_partitions}")
# Если данных 500 МБ, а партиций 50 000 - каждая по 10 КБ - проблема
approx_size_mb = 500
avg_partition_kb = (approx_size_mb * 1024) / num_partitions
print(f"Средний размер Spark-партиции: {avg_partition_kb:.1f} КБ")
if avg_partition_kb < 1024: # меньше 1 МБ
print("ПРОБЛЕМА: слишком мелкие партиции при чтении")
Диагностика через explain-план¶
df = spark.read.parquet("/path/to/table/")
query = df.filter(F.col("dt") == "2026-01-15").select("amount", "country")
query.explain(mode="formatted")
# Ищем в FileScan:
# Location: InMemoryFileIndex(50000 paths)[...] ← 50 000 файлов!
# numFiles: 50000
# sizeInBytes: 500.0 MiB
#
# Соотношение: 500 МБ / 50 000 файлов = 10 КБ / файл → проблема
В SQL/Details в Spark History Server ищем метрику FileScan:
- number of files read: если тысячи при небольшом Input Size - диагноз поставлен
- size of files read: суммарный объём. Делим на
number of files- получаем средний размер файла
5. Стратегии лечения¶
Лечение 1: Контроль числа файлов при записи¶
Самый простой способ предотвратить болезнь - контролировать число выходных файлов при каждой записи через coalesce или repartition.
import math
def write_optimally(df, path, target_file_mb=256, partition_col=None):
"""
Записать DataFrame с контролируемым числом выходных файлов.
target_file_mb: целевой размер файла (128-512 МБ - золотой стандарт).
"""
# Оцениваем размер Parquet через статистику плана (без full scan)
plan = df._jdf.queryExecution().optimizedPlan()
size_bytes = plan.stats().sizeInBytes()
# Parquet ≈ 30-50% от размера в памяти (эффект сжатия + колоночное хранение)
parquet_mb = float(size_bytes) / (1024 * 1024) * 0.4
num_files = max(1, math.ceil(parquet_mb / target_file_mb))
print(f"Оценка Parquet размера: ~{parquet_mb:.0f} МБ → {num_files} файлов")
writer = df.coalesce(num_files).write.mode("overwrite")
if partition_col:
writer.partitionBy(partition_col)
writer.parquet(path)
# Применение
events_df = spark.read.parquet("s3a://raw/events/dt=2026-01-15/")
write_optimally(events_df, "s3a://processed/events/dt=2026-01-15/", target_file_mb=256)
Когда coalesce vs repartition:
coalesce(n)- для уменьшения числа файлов без shuffle. Данные слипаются из существующих партиций. Быстро, но распределение может быть неравномерным. Используйте перед write если данные уже распределены нормально.repartition(n)- для равномерного перераспределения через shuffle. Нужен если данные skewed или если нужно увеличить число файлов. Медленнее, но выходные файлы одинаковые.repartition(n, col)- для записи сpartitionBy: создаёт ровноnфайлов в каждой partition-папке.
# Правило: groupBy создаёт 200 partition → 200 файлов на partition-папку
result = events_df \
.groupBy("event_date", "country") \
.agg(F.sum("amount").alias("total"), F.count("*").alias("cnt"))
# Решение: repartition по partition-ключу перед записью
result \
.repartition(2, "event_date", "country") \
.write \
.mode("overwrite") \
.partitionBy("event_date", "country") \
.parquet("s3a://output/daily_country_stats/")
# Результат: ровно 2 файла в каждой папке event_date=X/country=Y/
# Spark SQL-хинты для управления партициями
spark.sql("""
SELECT /*+ COALESCE(4) */ event_date, sum(amount) AS total
FROM events
WHERE event_date = '2026-01-15'
GROUP BY event_date
""").write.parquet("s3a://output/")
Лечение 2: AQE Partition Coalescing¶
AQE (Adaptive Query Execution) умеет автоматически объединять мелкие shuffle-партиции в более крупные после того, как увидит реальный объём данных. Это не решает проблему уже существующих мелких файлов в storage, но предотвращает создание новых при записи результатов.
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
# Целевой размер партиции после coalesce
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "268435456") # 256 МБ
# AQE будет стараться создавать партиции размером ~256 МБ
# Минимальное число партиций (не схлопывать в 1)
spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionNum", "1")
# Демонстрация без AQE:
spark.conf.set("spark.sql.adaptive.enabled", "false")
result_no_aqe = small_df.groupBy("category").agg(F.sum("amount"))
result_no_aqe.write.parquet("/tmp/without_aqe/")
# → 200 файлов по ~50 КБ (дефолт 200 shuffle partitions для маленького датасета)
# Демонстрация с AQE:
spark.conf.set("spark.sql.adaptive.enabled", "true")
result_with_aqe = small_df.groupBy("category").agg(F.sum("amount"))
result_with_aqe.write.parquet("/tmp/with_aqe/")
# AQE видит что данных мало → схлопывает 200 партиций в 1-3 файла оптимального размера
Лечение 3: spark.sql.files.maxPartitionBytes - объединение файлов при чтении¶
Этот параметр контролирует максимальный размер InputSplit при чтении. По умолчанию 128 МБ. Spark пытается объединить несколько мелких файлов в один InputSplit, если они суммарно меньше этого порога. Это не удаляет мелкие файлы с диска, но уменьшает число Tasks при чтении.
# По умолчанию: каждый файл = минимум 1 Task
# С maxPartitionBytes Spark может объединять файлы в один сплит
spark.conf.set("spark.sql.files.maxPartitionBytes", str(128 * 1024 * 1024)) # 128 МБ
# Если файлы по 10 МБ → Spark объединит ~13 файлов в один сплит → один Task
spark.conf.set("spark.sql.files.openCostInBytes", str(4 * 1024 * 1024)) # 4 МБ
# openCostInBytes - "цена" открытия файла в байтах данных.
# При 4 МБ: Spark считает открытие файла = 4 МБ накладных расходов.
# Эффект: файлы меньше 4 МБ объединяются агрессивнее - они "дорогие" для открытия.
# Практический эффект:
# До: 10 000 файлов по 10 МБ → 10 000 Tasks
# После: 10 000 × 10 МБ / 128 МБ ≈ 780 Tasks (в 12.8 раз меньше!)
Ограничение: maxPartitionBytes работает только для splittable форматов (Parquet без gzip-компрессии). Файлы .parquet.gz с gzip не объединяются этим механизмом - используйте Snappy или Zstd.
Лечение 4: Compaction - архитектурное решение¶
Compaction (уплотнение) - процесс слияния существующих мелких файлов в крупные. Это реактивное лечение для уже накопившихся файлов.
def compact_partition(table_path, partition_value, target_files=4):
"""Компакция конкретной партиции: читаем мелкие, пишем крупные."""
partition_path = f"{table_path}/event_date={partition_value}/"
df = spark.read.parquet(partition_path)
actual_files = df.rdd.getNumPartitions()
if actual_files <= target_files * 1.5:
print(f"Партиция {partition_value}: {actual_files} файлов - компакция не нужна")
return
print(f"Компакция: {partition_value}: {actual_files} → {target_files} файлов")
tmp_path = f"{partition_path}_compact_tmp/"
df.coalesce(target_files) \
.write \
.mode("overwrite") \
.parquet(tmp_path)
# Атомарная замена файлов
import shutil, os
shutil.rmtree(partition_path)
os.rename(tmp_path, partition_path)
print(f"Партиция {partition_value}: компакция завершена")
# Компакция последних 7 дней
from datetime import date, timedelta
today = date.today()
for days_ago in range(1, 8):
d = (today - timedelta(days=days_ago)).isoformat()
compact_partition("s3a://data/events", d, target_files=4)
Компакция через Delta Lake OPTIMIZE¶
Delta Lake предоставляет встроенную команду OPTIMIZE. Она:
- Читает все мелкие файлы в указанной партиции
- Объединяет их в файлы целевого размера (по умолчанию 1 ГБ)
- Атомарно заменяет старые файлы новыми через Transaction Log
- Старые файлы помечаются как deleted и физически удаляются через
VACUUM
# Оптимизировать всю таблицу
spark.sql("OPTIMIZE delta.events")
# Оптимизировать только конкретную партицию (эффективнее для больших таблиц)
spark.sql("""
OPTIMIZE delta.events
WHERE event_date = '2026-01-15'
""")
# OPTIMIZE + ZORDER: объединить И упорядочить данные по колонке
# Ускоряет последующие запросы с фильтрами по user_id
spark.sql("""
OPTIMIZE delta.events
ZORDER BY (user_id, country)
""")
# Посмотреть историю операций
spark.sql("DESCRIBE HISTORY delta.events").show(5, truncate=False)
# Удалить устаревшие файлы (минимум 7 дней retention для time-travel)
spark.sql("VACUUM delta.events RETAIN 168 HOURS")
Компакция через Apache Iceberg¶
В Iceberg компакция выполняется через Spark Procedure rewrite_data_files:
# Базовая компакция: объединить файлы меньше 64 МБ в файлы по 256 МБ
spark.sql("""
CALL iceberg.system.rewrite_data_files(
table => 'iceberg.events',
options => map(
'target-file-size-bytes', '268435456',
'min-file-size-bytes', '67108864'
)
)
""")
# Компакция с фильтром по партиции
spark.sql("""
CALL iceberg.system.rewrite_data_files(
table => 'iceberg.events',
where => 'event_date = ''2026-01-15''',
options => map('target-file-size-bytes', '268435456')
)
""")
# Компакция + сортировка (аналог ZORDER в Delta)
spark.sql("""
CALL iceberg.system.rewrite_data_files(
table => 'iceberg.events',
strategy => 'sort',
sort_order => 'user_id ASC NULLS LAST, event_date ASC',
options => map('target-file-size-bytes', '268435456')
)
""")
Лечение 5: Автоматическая компакция в Structured Streaming¶
Для streaming-пайплайнов правильная стратегия - встроить компакцию прямо в streaming-джобу.
# Delta Lake: Auto Optimize при записи из стрима
events_stream \
.writeStream \
.trigger(processingTime="1 minute") \
.format("delta") \
.option("path", "s3a://data/events_delta/") \
.option("checkpointLocation", "s3a://checkpoints/events/") \
.option("delta.autoOptimize.optimizeWrite", "true") \
# optimizeWrite: Delta сама решает сколько файлов создать в батче
# вместо N файлов (= числу Tasks) - минимальное нужное число
.option("delta.autoOptimize.autoCompact", "true") \
# autoCompact: после каждого батча Delta проверяет нужна ли компакция
# и запускает её в фоне если мелких файлов накопилось достаточно
.start()
# Альтернатива для plain Parquet: foreachBatch с явным контролем файлов
def write_batch_optimally(batch_df, batch_id):
"""Записать микробатч с контролем числа файлов + периодическая компакция."""
batch_df \
.repartition(4, "event_date") \
.write \
.mode("append") \
.partitionBy("event_date") \
.parquet("s3a://data/events/")
# Компакция каждые 100 батчей (≈ 1000 секунд при интервале 10с)
if batch_id % 100 == 0:
print(f"Batch {batch_id}: запуск плановой компакции...")
spark.read.parquet("s3a://data/events/") \
.coalesce(20) \
.write \
.mode("overwrite") \
.partitionBy("event_date") \
.parquet("s3a://data/events/")
events_stream \
.writeStream \
.trigger(processingTime="10 seconds") \
.foreachBatch(write_batch_optimally) \
.option("checkpointLocation", "s3a://checkpoints/events/") \
.start()
6. Профилактика: предотвратить лучше, чем лечить¶
Мониторинг числа файлов как операционная метрика¶
Число файлов в ключевых таблицах должно мониториться наравне с latency и error rate. Добавьте проверку в CI/CD или в ежедневные отчёты:
def check_table_health(path, warn_threshold=10_000, error_threshold=100_000):
"""Проверить здоровье таблицы по числу файлов."""
import os
file_count = 0
total_size = 0
for root, dirs, files in os.walk(path):
for f in files:
if f.endswith((".parquet", ".orc")):
fp = os.path.join(root, f)
file_count += 1
total_size += os.path.getsize(fp)
avg_size_mb = (total_size / file_count / (1024 * 1024)) if file_count > 0 else 0
status = "OK"
if file_count > warn_threshold:
status = "WARNING"
if file_count > error_threshold:
status = "CRITICAL"
print(f"[{status}] Файлов: {file_count:,} | Средний размер: {avg_size_mb:.1f} МБ")
if avg_size_mb < 32:
print(" Рекомендуется немедленная компакция")
elif avg_size_mb < 64:
print(" Рекомендуется плановая компакция")
return {"file_count": file_count, "avg_size_mb": avg_size_mb, "status": status}
# Запускать ежедневно для всех ключевых таблиц
check_table_health("s3a://data/events/")
check_table_health("s3a://data/orders/")
Чек-лист: «что проверить если джоба зависла на Listing»¶
Когда видите в логах:
INFO FileSourceScanExec: Listing leaf files and directories...
И это висит больше 30 секунд - Small Files Disease. Алгоритм диагностики:
# 1. Сколько файлов в таблице?
df = spark.read.parquet("/path/to/table/")
print(f"Spark-партиций (≈ файлов): {df.rdd.getNumPartitions()}")
# 2. Проверить explain-план
df.filter(F.col("dt") == "2026-01-15").explain(mode="formatted")
# Ищем: Location: InMemoryFileIndex(N paths) где N > 10 000
# 3. Сколько Tasks в Spark UI?
# Stages → конкретный Stage → смотрим Total Tasks
# Если > 10 000 Tasks для < 100 ГБ данных → проблема
# 4. Средняя длительность Task?
# В Spark UI Summary Metrics: если Median Duration < 100 мс → overhead доминирует
# 5. Применить быстрое лечение:
spark.conf.set("spark.sql.files.maxPartitionBytes", str(128 * 1024 * 1024))
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
# Это помогает при ЧТЕНИИ уже существующих мелких файлов
# Для долгосрочного решения нужна компакция
7. Сводная таблица и Best Practices¶
Сводная таблица техник борьбы с Small Files Disease¶
| Проблема | Решение | Когда применять |
|---|---|---|
| Много файлов при batch-записи | coalesce(N) или repartition(N) перед write |
После groupBy, join, filter |
| Мелкие shuffle-партиции при записи | AQE coalescePartitions.enabled = true |
Всегда при Spark 3.0+ |
| Мелкие файлы при чтении | spark.sql.files.maxPartitionBytes |
Быстрое облегчение без компакции |
| Накопившиеся файлы (Parquet) | Ручная компакция через coalesce + overwrite |
Разово или по расписанию |
| Накопившиеся файлы (Delta Lake) | OPTIMIZE / autoCompact |
Регулярно + после streaming |
| Накопившиеся файлы (Iceberg) | rewrite_data_files |
Регулярно по расписанию |
| Streaming без компакции | foreachBatch с repartition + периодическая компакция |
В каждой streaming-джобе |
| Streaming + Delta Lake | autoOptimize.optimizeWrite + autoCompact |
Всегда для Delta streaming |
Золотые правила¶
1. Целевой размер файла: 128–512 МБ. Перед каждой write-операцией думайте: «Сколько файлов создаст этот код?» Используйте coalesce(N) или repartition(N) для контроля.
2. Включить AQE. adaptive.enabled = true + coalescePartitions.enabled = true - это должно быть в базовой конфигурации каждого кластера.
3. Настроить maxPartitionBytes. 128 МБ по умолчанию подходит для большинства случаев. При очень мелких файлах можно попробовать 256 МБ.
4. Регулярная компакция для streaming-таблиц. Раз в день или раз в неделю - зависит от интенсивности потока. В Delta Lake autoCompact делает это автоматически.
5. Мониторинг числа файлов. Добавьте check_table_health() в ежедневные healthcheck-джобы.
Домашнее задание¶
Дан следующий streaming-пайплайн:
api_logs_stream \
.writeStream \
.trigger(processingTime="30 seconds") \
.partitionBy("request_date", "endpoint_category") \
.format("parquet") \
.option("path", "s3a://data/api_logs/") \
.start()
# endpoint_category: 500 уникальных значений
# request_date: обновляется ежедневно
# Объём данных: 1 ГБ/час
Задание 1 - Посчитать катастрофу. Рассчитайте, сколько файлов накопится за 30 дней: batches в день × файлы на batch × число partition-значений. При processingTime = 30 seconds и 500 уникальных endpoint_category.
Задание 2 - Предложить архитектуру. Переработайте streaming write: используйте foreachBatch, контролируйте число файлов через repartition(), добавьте периодическую компакцию каждые N батчей. Обоснуйте выбор N.
Задание 3 - Delta Lake решение. Перепишите пайплайн с Delta Lake + autoOptimize.autoCompact. Добавьте OPTIMIZE джоб по расписанию. Объясните почему Delta решает проблему лучше чем plain Parquet в этом сценарии.
Задание 4 - Функция диагностики. Напишите функцию, которая принимает путь к таблице и возвращает: число файлов, средний размер, min/max/avg по партициям, рекомендации (компакция нужна / не нужна / критично). Протестируйте на синтетических данных с намеренно созданными мелкими файлами.
В следующем уроке разберём стратегии компакции подробнее: как настроить автоматическую компакцию в Delta Lake и Iceberg, как выбрать расписание, как OPTIMIZE взаимодействует с Z-Order и Liquid Clustering, и как не сломать time-travel при агрессивной чистке через VACUUM.