Small Files Disease: причины, последствия и влияние на планировщик

Почему 1000 файлов по 1 КБ медленнее одного файла в 1 МБ, как Driver падает с OOM от метаданных, откуда берётся болезнь в streaming и CDC пайплайнах, и как лечить через Compaction, AQE и современные Lakehouse форматы.

optimization internals

Есть проблема, которая встречается в каждом 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 обязан:

  1. Собрать список всех файлов для чтения (Directory Listing)
  2. Определить, какие partition-папки нужны (Partition Pruning)
  3. Разбить файлы на InputSplits для Tasks
  4. Создать объект TaskDescription для каждого Task
  5. Сериализовать и отправить 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.