Small Files на S3: LIST latency и ограничения при большом числе файлов

Small Files Problem - главный инфраструктурный капкан при переходе на S3. Разбираем физику HTTP-накладных расходов, Scheduler Overhead Spark, стратегии компакции и как Lakehouse-форматы решают проблему метаданных.

storage

Почему Small Files - это архитектурная угроза, а не косметика

Когда инженер впервые слышит «проблема мелких файлов», реакция обычно скептическая: «Ну файлы маленькие, зато их много. Данных-то столько же». Это опасное заблуждение.

Проблема Small Files в облачном объектном хранилище - это не просто «чуть медленнее». Это системная деградация, которая проявляется сразу на нескольких уровнях:

  • Уровень S3 API: 1 миллион файлов = 1 миллион HTTP GET-запросов при чтении. Каждый запрос - 5–50 мс латентности TLS-соединения. 1M × 20 мс = 5,5 часов чистого I/O Wait.
  • Уровень метаданных: LIST операция возвращает максимум 1000 объектов за раз. Для листинга 1М файлов - 1000 запросов ListObjectsV2 только чтобы «понять, что читать».
  • Уровень Spark Planner: каждый файл - потенциально отдельный таск. 1М файлов = 1М тасок = OutOfMemory на Spark Driver.
  • Уровень денег: каждый API-вызов к S3 тарифицируется. 1М GET запросов = $0.40. В год при ежедневных запросах - $146. При 100 Spark-джобах в сутки - $14 600 в год только за API-вызовы к мелким файлам.

Проблема не появляется «внезапно». Она нарастает постепенно: сначала запросы замедляются с 2 до 5 минут, потом до 30 минут, потом джобы начинают падать по timeout. К моменту когда это становится «серьёзной проблемой» - в озере уже десятки миллионов файлов, и исправить это без downtime очень сложно.

Этот урок - про физику проблемы, её источники и способы предотвращения и лечения.

Физика проблемы: почему S3 - это не жёсткий диск

Объектное хранилище - это веб-сервис

HDFS на локальных дисках: чтение файла = системный вызов read() к локальной файловой системе. Задержка 0.1–1 мс. Последовательное чтение блока = несколько мкс на блок.

S3: чтение файла = HTTP GET запрос к удалённому веб-сервису. Каждый запрос включает:

  1. DNS-резолюцию (если соединение новое): ~5 мс.
  2. TCP-хэндшейк: ~1–2 RTT.
  3. TLS-хэндшейк: ~1–2 RTT дополнительно.
  4. Отправку HTTP-запроса.
  5. Ожидание первого байта ответа (TTFB): 5–50 мс.
  6. Получение данных по сети.

Time To First Byte (TTFB) для S3 - это «плата за вход» в каждый файл. Независимо от размера файла - 1 байт или 1 ГБ - TTFB одинаков.

Математика TTFB:

Случай 1: 1 файл 128 MB
- 1 HTTP GET запрос
- TTFB: 20 мс
- Передача данных: 128 MB / 100 MB/s = 1280 мс
- Итого: 1300 мс, КПД = 1280/1300 = 98.5%

Случай 2: 10 000 файлов по 12.8 KB (итого тоже 128 MB)
- 10 000 HTTP GET запросов
- TTFB суммарно: 10 000 × 20 мс = 200 000 мс = 3.3 МИНУТЫ
- Передача данных: 128 MB / 100 MB/s = 1280 мс
- Итого: 201 280 мс = 3.4 МИНУТЫ
- КПД = 1280 / 201280 = 0.6% (!!)

Ускорение от перехода к крупным файлам: 201280 / 1300 = 155×

При параллельном выполнении (100 воркеров одновременно) картина улучшается, но не принципиально: S3 начинает throttle при высоком RPS.

Тарификация S3 API: скрытые расходы

AWS тарифицирует каждый API-вызов к S3:

Операция Стоимость
PUT, COPY, POST, LIST $0.005 за 1000 запросов
GET, SELECT $0.0004 за 1000 запросов
Хранение $0.023 за ГБ/месяц

Расчёт для типичного Data Lake:

Сценарий: 100 Spark-джобов в сутки, каждый читает таблицу из 50 000 файлов

Запросов в сутки: 100 × 50 000 = 5 000 000 GET
Стоимость GET: 5 000 000 / 1000 × $0.0004 = $2.00 / день
В месяц: $60
В год: $730

Плюс LIST при планировании (50 000 / 1000 = 50 ListObjectsV2 per job):
100 × 50 = 5 000 LIST в сутки
5 000 / 1000 × $0.005 = $0.025 / день → $9 / год

Итого только API: ~$730 / год за один Data Lake

Тот же объём данных в 100 файлах:
100 × 100 = 10 000 GET / сутки
10 000 / 1000 × $0.0004 = $0.004 / день → $1.46 / год

Экономия: $728 / год только за API-вызовы

При масштабе крупного Data Lake с тысячами джобов и миллиардами файлов - это десятки тысяч долларов в год только за API.

LIST Latency: метаданные как бутылочное горлышко

Как S3 листает директории

В S3 нет настоящих директорий (это мы разбирали в уроке про Rename Problem). «Директория» - это просто общий prefix ключей. Операция ListObjectsV2 с указанным prefix возвращает все объекты, начинающиеся с этого prefix.

Ограничение: S3 возвращает максимум 1000 объектов за один вызов ListObjectsV2. Для листинга большего количества нужна пагинация.

Это для 5000 файлов. Для 1 000 000 файлов - 1000 запросов ListObjectsV2, каждый с TTFB ~20 мс = 20 секунд только на листинг. А Spark ещё не начал читать данные!

Рекурсивный листинг партиционированных таблиц

Партиционированная таблица - ещё хуже. Spark должен залистить каждую директорию отдельно:

events/
├── event_date=2024-01-01/  ← ListObjectsV2 №1
│   ├── hour=00/            ← ListObjectsV2 №2
│   ├── hour=01/            ← ListObjectsV2 №3
│   └── ...hour=23/         ← ListObjectsV2 №25
├── event_date=2024-01-02/  ← ListObjectsV2 №26
│   ├── hour=00/            ← ListObjectsV2 №27
│   └── ...
└── event_date=2024-12-31/  ← ListObjectsV2 №N

365 дней × 24 часа = 8760 листингов только для обнаружения структуры таблицы!
Если в каждой директории 100 файлов → 8760 ListObjectsV2 запросов
8760 × 20мс = 175 секунд только на планирование

Throttling S3: когда много запросов одновременно

S3 лимитирует количество запросов на prefix (partition key внутри S3):

  • 5500 GET/HEAD запросов в секунду на prefix.
  • 3500 PUT/COPY/POST/DELETE запросов в секунду на prefix.

Когда Spark запускает 200 воркеров, каждый параллельно листит свои файлы:

200 воркеров × 50 GET/сек = 10 000 GET/сек > 5500 лимит
→ HTTP 503 SlowDown
→ Spark получает ошибки, повторяет запросы
→ Повторные запросы создают ещё большую нагрузку
→ Каскадный throttling → джоб зависает или падает по таймауту

Откуда берутся мелкие файлы

Источник 1: Streaming micro-batch

Каждый micro-batch Structured Streaming создаёт новые файлы:

# ❌ Каждые 30 секунд - новые файлы:
df.writeStream \
    .trigger(processingTime="30 seconds") \
    .format("parquet") \
    .option("path", "s3a://bucket/events/") \
    .start()

# Расчёт: 30 секунд / trigger × 86400 секунд/день = 2880 batch в сутки
# Если каждый batch создаёт 10 файлов → 28 800 файлов в сутки
# За месяц: 864 000 файлов
# За год: 10 368 000 файлов

Источник 2: Высококардинальное партиционирование

# ❌ Партиционирование по user_id - катастрофа!
df.write \
    .partitionBy("user_id") \    # 10 миллионов уникальных user_id
    .parquet("s3a://bucket/events/")

# Результат: 10 000 000 директорий, в каждой - один крошечный файл
# Spark создаёт МИНИМУМ по одному файлу на партицию × на воркер!
# 10M user_id × 200 воркеров = потенциально 2 000 000 000 файлов

Источник 3: Чрезмерный repartition перед записью

# ❌ "Для параллелизма" установим 1000 партиций:
df \
    .repartition(1000) \    # 1000 shuffle партиций
    .write \
    .parquet("s3a://bucket/output/")

# Если датасет = 500 MB, то каждый файл = 500 KB
# 1000 файлов по 500 KB вместо 4 файлов по 128 MB

Источник 4: Incremental ETL без контроля размера

# ❌ Ежечасный incremental load без coalesce:
hourly_events = spark.read.kafka(...) \
    .filter(f"hour = '{current_hour}'")

hourly_events.write \
    .mode("append") \
    .partitionBy("date", "hour") \        # каждый час - новые файлы
    .parquet("s3a://bucket/events/")

# Если данных за час мало (ночью, в выходные):
# 50 MB данных / 128 MB по умолчанию → 1 файл, но...
# Spark создаёт файл на каждого воркера, который что-то писал
# 50 активных воркеров → 50 файлов по 1 MB

Как Spark страдает от мелких файлов: три уровня деградации

Уровень 1: Planning Overhead - Driver OutOfMemory

Spark Driver перед запуском задачи должен «открыть» все входные файлы и построить план:

  1. Рекурсивный листинг директорий (ListObjectsV2).
  2. Получение FileStatus для каждого файла (size, modification_time).
  3. Разбивка файлов на InputSplits (один Split = один таск по умолчанию).
  4. Сериализация InputSplits и отправка воркерам.

Для 1 000 000 файлов:

  • Список FileStatus: 1M объектов × ~200 bytes каждый = 200 MB в памяти Driver.
  • Список InputSplits: 1M объектов × ~500 bytes = 500 MB.
  • Итого только на метаданные: ~700 MB в памяти Driver.
  • spark.driver.memory = 4g? Останется 3.3 GB для всего остального.
  • spark.driver.memory = 2g? Driver падает с OutOfMemoryError.
# Симптомы в логах:
# java.lang.OutOfMemoryError: Java heap space
# at org.apache.hadoop.fs.s3a.S3AFileSystem.listFiles(S3AFileSystem.java:...)
# при попытке планирования задачи над 1M+ файлами

# Диагностика: посмотреть сколько файлов будет обработано:
df = spark.read.parquet("s3a://bucket/events/")
print(f"Количество входных файлов: {len(df.inputFiles())}")
print(f"Количество партиций: {df.rdd.getNumPartitions()}")

# ❌ Если вывод: 850 000 файлов, 850 000 партиций → критическая проблема

Уровень 2: Task Scheduling Overhead

Spark планировщик управляет тасками через структуры данных в памяти Driver:

1 000 000 тасок × состояние таска (PENDING/RUNNING/SUCCEEDED):
- TaskDescription: ~1 KB каждый → 1 GB только на описание тасков
- TaskResult: метрики каждого завершённого таска → ещё несколько ГБ

Накладные расходы планировщика:
- Назначение 1M тасок на 200 воркеров: O(N × executors) = O(200M) операций
- Heartbeat от каждого воркера: 200 × каждые 10 сек = 20 heartbeat/сек
- Сериализация TaskDescription для отправки: 1 KB × 1M = 1 GB сетевого трафика
  только для раздачи тасок (не данных!)

Для 1M тасок по 10 мс работы каждая (чтение 10KB файла):

  • Чистая работа: 1M × 10 мс / 200 воркеров = 50 секунд вычислений.
  • Scheduling overhead: 1M тасок × 5 мс overhead = 83 минуты.
  • Итого: 84 минуты вместо 50 секунд. КПД = 1%.

Уровень 3: Vectorized Reader Degradation

Как мы разбирали в уроке о Parquet - vectorized reader читает данные батчами по 4096 строк. Для файла 10 KB:

  • 10 KB / 50 bytes/row ≈ 200 строк.
  • Одного batcha (4096 строк) не набирается - файл меньше одного batch!
  • Vectorized Reader переключается в деградированный режим.
  • Footer читается отдельным HTTP GET запросом → это 50% объёма данных!
Оверхед Footer для мелких файлов:

Файл 1 GB: Footer обычно 100 KB = 0.01% overhead - незаметно
Файл 128 MB: Footer 50 KB = 0.04% overhead - незаметно
Файл 1 MB: Footer 20 KB = 2% overhead - уже заметно
Файл 10 KB: Footer 5 KB = 50% overhead - половина запроса - на метаданные!
Файл 1 KB: Footer 2 KB = 200% overhead - больше читаем метаданных чем данных

Spark-механизмы борьбы с мелкими файлами

Автоматическое склеивание: maxPartitionBytes и openCostInBytes

Spark пытается объединять мелкие файлы в одну задачу через параметры:

# Целевой размер одной Spark-партиции (таска):
spark.conf.set("spark.sql.files.maxPartitionBytes", str(128 * 1024 * 1024))  # 128 MB

# «Штраф» за открытие каждого нового файла (учитывается как добавка к размеру):
spark.conf.set("spark.sql.files.openCostInBytes", str(4 * 1024 * 1024))  # 4 MB

# Как работает:
# Файлы: [5MB, 3MB, 4MB, 6MB, 8MB, ...]
# Spark объединяет их в один таск пока сумма < maxPartitionBytes
# openCostInBytes добавляется к каждому файлу: фактический "вес" = 3MB + 4MB = 7MB
# Это делает маленькие файлы "тяжелее" → меньше объединяется в один таск
# → больше параллелизма при чтении крупных и мелких файлов вместе

Почему это не решает проблему полностью:

Даже если Spark объединяет 100 файлов по 1 MB в один таск - он всё равно делает 100 отдельных HTTP GET запросов к S3. Параллельно, да, но это 100 TTFB = 100 × 20 мс = 2 секунды накладных расходов только на «открытие» данных для одного таска. Против 20 мс для одного файла 100 MB.

Adaptive Query Execution (AQE): динамическое объединение

Spark 3.x добавил AQE - адаптивную оптимизацию во время выполнения:

# Включить AQE:
spark.conf.set("spark.sql.adaptive.enabled", "true")

# Автоматическое объединение мелких партиций после shuffle:
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionSize",
               str(1 * 1024 * 1024))    # минимальный размер партиции: 1 MB
spark.conf.set("spark.sql.adaptive.coalescePartitions.initialPartitionNum",
               "1000")                  # начальное число партиций shuffle

# Как работает AQE coalescePartitions:
# После shuffle Spark видит 500 партиций по 100 KB каждая (результат агрегации)
# AQE объединяет их в ≈ ceil(50MB / 128MB) = 1 партицию
# Вместо 500 тасок на следующий stage → 1 таск
# НО: это применяется только к shuffle output, не к входным файлам!

Важное ограничение: AQE помогает с мелкими партициями после shuffle операций, но не решает проблему мелких файлов при чтении с S3. Листинг всё равно происходит, HTTP-запросы всё равно делаются.

Стратегии решения: проактивная и реактивная

Стратегия 1: Правильный coalesce/repartition перед записью

from pyspark.sql import functions as F
import math


def write_with_optimal_file_count(
    df,
    output_path: str,
    target_file_size_mb: int = 256,
    partition_col: str = None
) -> None:
    """
    Записать DataFrame с оптимальным количеством файлов.

    Логика: оцениваем размер данных, вычисляем нужное число файлов.
    coalesce() - более эффективен (нет shuffle), но менее равномерен.
    repartition() - создаёт равные части через shuffle.
    """
    # Оценить размер данных в байтах:
    # (приблизительно через plan stats если доступно)
    try:
        plan = df._jdf.queryExecution().analyzed()
        estimated_bytes = plan.stats().sizeInBytes()
        target_files = max(1, math.ceil(estimated_bytes / (target_file_size_mb * 1024 * 1024)))
    except Exception:
        # Фоллбэк: использовать число текущих партиций
        current_parts = df.rdd.getNumPartitions()
        target_files = max(1, current_parts // 4)  # объединить в 4× меньше

    print(f"Целевое количество файлов: {target_files}")

    writer = (
        df
        .coalesce(target_files)       # coalesce уменьшает без shuffle
        .write
        .option("parquet.block.size", target_file_size_mb * 1024 * 1024)
        .mode("overwrite")
    )

    if partition_col:
        writer = writer.partitionBy(partition_col)

    writer.parquet(output_path)


# Правило: 1 файл на воркер для batch задач,
# но не меньше чем spark.default.parallelism / 4

Стратегия 2: Контроль partitionBy с низкой кардинальностью

# ❌ Высококардинальный partition → миллионы директорий:
df.write.partitionBy("user_id").parquet(path)  # 10M уникальных = 10M директорий

# ✅ Партиционировать по низкокардинальным полям:
df.write.partitionBy("event_date", "country").parquet(path)
# 365 × 200 стран = максимум 73 000 директорий (реально меньше)

# ✅ Bucket partition для высококардинальных полей:
# Разбить user_id на N бакетов:
df.withColumn("user_bucket", F.pmod(F.hash("user_id"), F.lit(256))) \
    .write \
    .partitionBy("event_date", "user_bucket") \  # 365 × 256 = 93 440 директорий max
    .parquet(path)
# Запрос WHERE user_id = 12345:
# → hash(12345) % 256 = 73 → читаем только bucket 73
# Row Group Filtering внутри bucket по user_id

Стратегия 3: Streaming с контролем числа файлов

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


def start_streaming_with_file_control(
    spark: SparkSession,
    source_df,
    output_path: str,
    target_files_per_trigger: int = 4,
    trigger_interval: str = "5 minutes"
):
    """
    Structured Streaming с контролем количества файлов.

    target_files_per_trigger: желаемое число файлов на trigger.
    Используем foreachBatch для контроля coalesce.
    """
    def write_batch(batch_df, batch_id: int):
        if batch_df.isEmpty():
            return

        batch_df \
            .coalesce(target_files_per_trigger) \
            .write \
            .mode("append") \
            .partitionBy("event_date") \
            .option("parquet.block.size", str(128 * 1024 * 1024)) \
            .parquet(output_path)

    return (
        source_df
        .writeStream
        .foreachBatch(write_batch)
        .trigger(processingTime=trigger_interval)
        .option("checkpointLocation", f"{output_path}/_checkpoint/")
        .start()
    )


# Вместо 2880 batch/сутки × 50 файлов = 144 000 файлов/сутки
# Получаем: 2880 batch × 4 файла = 11 520 файлов/сутки
# Но: лучше увеличить интервал trigger!

# ✅ Оптимальный вариант: trigger раз в 15-30 минут:
# 96 trigger/сутки × 4 файла = 384 файла/сутки

Стратегия 4: Compaction Pipeline (реактивный подход)

Если мелкие файлы уже накопились - нужна компакция. Это периодически запускаемый Spark-джоб, который «склеивает» файлы:

from pyspark.sql import SparkSession, DataFrame
from pyspark.sql import functions as F
from dataclasses import dataclass
from datetime import datetime, timedelta
from typing import Optional
import logging

logger = logging.getLogger(__name__)


@dataclass
class CompactionConfig:
    """Конфигурация пайплайна компакции."""
    source_path: str
    output_path: str             # может совпадать с source (overwrite)
    partition_col: str           # колонка партиционирования для выбора "вчера"
    target_file_size_mb: int = 256
    sort_cols: list = None       # сортировка для оптимального Predicate Pushdown
    lookback_days: int = 1       # сколько дней компактировать


def run_compaction(
    spark: SparkSession,
    cfg: CompactionConfig,
    date_to_compact: Optional[str] = None
) -> dict:
    """
    Компакция партиции: читает мелкие файлы за дату,
    склеивает в крупные, перезаписывает партицию.

    Идемпотентно: повторный запуск безопасен.
    """
    if date_to_compact is None:
        # По умолчанию: компактировать вчерашние данные
        yesterday = (datetime.now() - timedelta(days=1)).strftime("%Y-%m-%d")
        date_to_compact = yesterday

    logger.info(f"Компакция {cfg.source_path} за {date_to_compact}")

    # Читаем только нужную партицию:
    partition_filter = f"{cfg.partition_col} = '{date_to_compact}'"
    df = spark.read.parquet(cfg.source_path).filter(partition_filter)

    # Статистика ПЕРЕД компакцией:
    input_files = df.inputFiles()
    file_count_before = len(input_files)
    logger.info(f"Файлов до компакции: {file_count_before}")

    if file_count_before <= 4:
        logger.info("Файлов уже мало, компакция не нужна")
        return {"status": "skipped", "files_before": file_count_before}

    # Вычислить целевое число файлов:
    # Spark Plan Stats (приблизительный размер):
    size_estimate = df._jdf.queryExecution().analyzed().stats().sizeInBytes()
    target_files = max(1, int(size_estimate / (cfg.target_file_size_mb * 1024 * 1024)) + 1)
    logger.info(f"Целевое число файлов: {target_files} "
                f"(оценка размера: {size_estimate / 1024**2:.1f} MB)")

    # Сортировка для Predicate Pushdown:
    if cfg.sort_cols:
        df = df.sortWithinPartitions(*cfg.sort_cols)

    # Запись в оптимизированном формате:
    (
        df
        .coalesce(target_files)
        .write
        .mode("overwrite")
        .partitionBy(cfg.partition_col)
        .option("parquet.block.size", cfg.target_file_size_mb * 1024 * 1024)
        .parquet(cfg.output_path)
    )

    logger.info(f"Компакция завершена: {file_count_before}{target_files} файлов")

    return {
        "status": "compacted",
        "date": date_to_compact,
        "files_before": file_count_before,
        "files_after": target_files,
        "reduction_ratio": round(file_count_before / target_files, 1),
    }


# Запуск из Airflow DAG:
# def compaction_task(**kwargs):
#     spark = create_spark_session()
#     cfg = CompactionConfig(
#         source_path="s3a://data-lake/events/",
#         output_path="s3a://data-lake/events/",
#         partition_col="event_date",
#         target_file_size_mb=256,
#         sort_cols=["user_id", "event_type"]
#     )
#     result = run_compaction(spark, cfg)
#     return result

Стратегия 5: Delta Lake OPTIMIZE - встроенная компакция

# Delta Lake имеет встроенную команду OPTIMIZE:
spark.sql("""
    OPTIMIZE delta.`s3a://data-lake/events`
    WHERE event_date = '2024-06-15'
""")

# OPTIMIZE:
# 1. Находит все файлы в партиции
# 2. Объединяет их в файлы размером ~1 GB (по умолчанию)
# 3. Атомарно заменяет через Delta Log
# 4. Старые мелкие файлы помечаются удалёнными (но хранятся до VACUUM)

# С Z-ordering для оптимального Predicate Pushdown:
spark.sql("""
    OPTIMIZE delta.`s3a://data-lake/events`
    ZORDER BY (user_id, event_date)
""")

# Настройка целевого размера файла:
spark.conf.set("spark.databricks.delta.optimize.maxFileSize",
               str(256 * 1024 * 1024))  # 256 MB (Databricks)
# или через table properties:
spark.sql("""
    ALTER TABLE delta.`s3a://data-lake/events`
    SET TBLPROPERTIES ('delta.targetFileSize' = '268435456')
""")

# После OPTIMIZE - VACUUM для удаления старых мелких файлов:
spark.sql("""
    VACUUM delta.`s3a://data-lake/events`
    RETAIN 168 HOURS   -- хранить историю 7 дней для time travel
""")

Стратегия 6: Apache Iceberg rewrite_data_files

# Iceberg аналог OPTIMIZE:
spark.sql("""
    CALL spark_catalog.system.rewrite_data_files(
        table => 'catalog.events',
        strategy => 'binpack',
        options => map(
            'target-file-size-bytes', '268435456',    -- 256 MB
            'min-file-size-bytes', '134217728',       -- компактировать файлы < 128 MB
            'max-concurrent-file-group-rewrites', '5' -- параллелизм
        ),
        where => 'event_date = ''2024-06-15'''
    )
""")

# С Sort Order для Predicate Pushdown:
spark.sql("""
    CALL spark_catalog.system.rewrite_data_files(
        table => 'catalog.events',
        strategy => 'sort',
        sort_order => 'user_id ASC NULLS LAST, event_date ASC NULLS LAST',
        options => map('target-file-size-bytes', '268435456')
    )
""")

# Результат метаданных (без физического копирования):
spark.sql("""
    CALL spark_catalog.system.expire_snapshots(
        table => 'catalog.events',
        older_than => TIMESTAMP '2024-06-08 00:00:00',
        retain_last => 10
    )
""")

Как Lakehouse-форматы решают проблему метаданных

Принцип: замена S3 LIST на чтение метаданных таблицы

Главная причина, почему мелкие файлы так дорого стоят на S3 - это дорогой рекурсивный LIST при планировании Spark-запроса. Delta Lake и Iceberg полностью исключают S3 LIST из критического пути:

Ключевое преимущество: при 1 000 000 файлов:

  • Plain Parquet: 1000 ListObjectsV2 запросов + обработка 1M FileStatus.
  • Delta Lake: 1–2 GET запроса к _delta_log/ + фильтрация в памяти по JSON.
  • Iceberg: 3–5 GET запросов к manifest + фильтрация.

Это разница в 200–1000 раз на этапе планирования.

Metadata Checkpointing в Delta Lake

Delta Log хранит транзакции в JSON-файлах. После 10 транзакций создаётся checkpoint - Parquet-файл с суммарным состоянием:

# Посмотреть структуру Delta Log:
import os

delta_log_path = "s3a://data-lake/events/_delta_log/"

# Файлы вида:
# 00000000000000000000.json   ← первая транзакция
# 00000000000000000001.json
# ...
# 00000000000000000009.json
# 00000000000000000010.checkpoint.parquet  ← checkpoint после 10 транзакций
# 00000000000000000011.json
# 00000000000000000012.json

# При запросе: Spark читает последний checkpoint + JSON после него
# Для таблицы с 1000 транзакций: 1 checkpoint.parquet + до 10 JSON
# Независимо от числа файлов в таблице!

# Настроить частоту checkpoint:
spark.conf.set("spark.databricks.delta.checkpointInterval", "10")  # каждые 10 коммитов

Manifest Pruning в Iceberg

Iceberg's manifest files содержат per-file column statistics. При большом количестве файлов манифесты тоже разбиваются на группы:

# Структура Iceberg metadata tree:
# metadata/v1.metadata.json
#   └── snap-12345.avro  (Manifest List)
#       ├── manifest-000.avro  (описывает файлы в partition 2024-01)
#       │   ├── data/part-001.parquet  min=2024-01-01, max=2024-01-31
#       │   └── data/part-002.parquet  min=2024-01-05, max=2024-01-31
#       ├── manifest-001.avro  (описывает файлы в partition 2024-02)
#       └── manifest-002.avro  (описывает файлы в partition 2024-03)

# При запросе WHERE event_date = '2024-02-15':
# 1. Читаем Manifest List (1 GET) → видим, что только manifest-001 может содержать 2024-02
# 2. Читаем manifest-001.avro (1 GET) → список файлов с min/max
# 3. Фильтруем файлы по min/max
# 4. Читаем только нужные файлы данных

# Из 1000 манифест-файлов читаем 1!

# Настроить параметры manifest:
spark.sql("""
    ALTER TABLE catalog.events SET TBLPROPERTIES (
        'write.metadata.metrics.mode' = 'full',
        'write.parquet.row-group-size-bytes' = '268435456'
    )
""")

Диагностика: аудит Small Files Problem

Подсчёт файлов и оценка ситуации

from pyspark.sql import SparkSession
import boto3
from collections import defaultdict


def audit_s3_files(
    spark: SparkSession,
    path: str,
    size_thresholds_mb: list = None
) -> dict:
    """
    Аудит файлов в S3: сколько файлов, их распределение по размеру.
    """
    if size_thresholds_mb is None:
        size_thresholds_mb = [1, 10, 64, 128, 256, 512]

    # Через Spark:
    df = spark.read.format("binaryFile").load(path)
    file_stats = df.select("path", "length").collect()

    total_files = len(file_stats)
    total_size_mb = sum(f.length for f in file_stats) / 1024 / 1024

    # Распределение по размеру:
    buckets = defaultdict(int)
    for f in file_stats:
        size_mb = f.length / 1024 / 1024
        for threshold in size_thresholds_mb:
            if size_mb < threshold:
                buckets[f"< {threshold} MB"] += 1
                break
        else:
            buckets[f"> {size_thresholds_mb[-1]} MB"] += 1

    avg_size_mb = total_size_mb / total_files if total_files else 0

    result = {
        "total_files": total_files,
        "total_size_mb": round(total_size_mb, 2),
        "avg_size_mb": round(avg_size_mb, 2),
        "size_distribution": dict(buckets),
        "severity": (
            "CRITICAL" if total_files > 100_000 else
            "HIGH" if total_files > 10_000 else
            "MEDIUM" if total_files > 1_000 else
            "OK"
        )
    }

    # Расчёт стоимости API:
    # LIST: total_files / 1000 × $0.005
    list_cost_per_query = (total_files / 1000) * 0.005
    # GET: total_files × $0.0004 / 1000
    get_cost_per_query = total_files * 0.0004 / 1000

    result["api_cost_per_full_scan"] = {
        "list_usd": round(list_cost_per_query, 4),
        "get_usd": round(get_cost_per_query, 4),
        "total_usd": round(list_cost_per_query + get_cost_per_query, 4),
    }

    return result


# Использование:
audit = audit_s3_files(spark, "s3a://data-lake/events/")
print(f"Файлов: {audit['total_files']}")
print(f"Средний размер: {audit['avg_size_mb']} MB")
print(f"Серьёзность: {audit['severity']}")
print(f"API стоимость за одно полное сканирование: ${audit['api_cost_per_full_scan']['total_usd']}")

Spark UI: признаки Small Files Problem

В Spark UI ищем аномальные показатели:

Вкладка Stages → конкретный stage:

1. Число тасок непропорционально большое:
   - 50 000 тасок при обработке 500 MB данных → Small Files Problem
   - Нормально: 500 MB / 128 MB = 4 таска

2. Task Deserialization Time высокое:
   - Если Deserialization > Compute → overhead на сериализацию TaskDescription
   - Симптом: Driver тратит больше времени на управление тасками, чем воркеры на работу

3. Scheduler Delay высокое:
   - Scheduler Delay > 100 мс → очередь на назначение тасков переполнена
   - Нормально: Scheduler Delay < 10 мс

4. GC Time высокое на Driver:
   - Driver собирает FileStatus 1M файлов → GC давление
   - GC Time > 10% от общего времени → проблема

Вкладка SQL → FileScan оператор:
- "number of files read" = 50 000+ → проблема
- "metadata time" > "scan time" → листинг дороже чтения данных

Anti-patterns: типичные ошибки и как их исправить

Антипаттерн 1: coalesce(1) для «одного файла»

# ❌ Антипаттерн: принудительно один файл
df.coalesce(1).write.parquet(output_path)

# Проблемы:
# 1. Всё данные проходят через один воркер → нет параллелизма → медленно
# 2. Файл может быть огромным (500 GB в одном файле?) → Spark не может его разделить
# 3. При чтении: один файл = один таск = нет параллелизма

# ✅ Правильно: оптимальное число файлов
num_partitions = max(1, spark.sparkContext.defaultParallelism)
df.coalesce(num_partitions).write.parquet(output_path)
# или для больших данных:
# df.repartition(num_partitions).write.parquet(output_path)

Антипаттерн 2: Hourly partitioning с низким объёмом данных

# ❌ Партиционирование по часу при низком объёме:
df.write \
    .partitionBy("year", "month", "day", "hour") \  # 8760 директорий в год
    .parquet(path)

# Если в час приходит 5 MB данных:
# 5 MB × 8760 часов = 43 GB / год (нормально)
# НО: 8760 директорий × 50 файлов в каждой = 438 000 файлов

# ✅ Партиционировать по дню, внутри - сортировать по часу:
df.withColumn("event_hour", F.hour("event_timestamp")) \
    .sortWithinPartitions("event_hour") \
    .write \
    .partitionBy("event_date") \        # только 365 директорий в год
    .parquet(path)

# Фильтр по часу работает через Row Group Filtering:
df_read.filter(
    (F.col("event_date") == "2024-06-15") &
    (F.col("event_hour") == 14)  # час как данные, не партиция
)

Антипаттерн 3: Бесконечный append без компакции

# ❌ Бесконечный append без очистки:
for i in range(365):  # 365 дней
    daily_data = load_daily_data(i)
    daily_data.write \
        .mode("append") \                # каждый день добавляет файлы
        .partitionBy("date") \
        .parquet(path)

# Через год: 365 × 50 воркеров × 1 файл = 18 250 файлов
# Через 5 лет: 91 250 файлов

# ✅ Компакция по расписанию + overwrite:
def incremental_with_compaction(daily_data, path: str, date: str):
    # Сначала записываем append:
    daily_data.write \
        .mode("append") \
        .partitionBy("event_date") \
        .parquet(path)

    # Потом компактируем эту партицию:
    partition_path = f"{path}/event_date={date}/"
    partition_data = spark.read.parquet(partition_path)
    partition_data \
        .coalesce(4) \               # оптимальное число файлов
        .sortWithinPartitions("user_id") \
        .write \
        .mode("overwrite") \
        .parquet(partition_path)

Production-кейс: деградация от 2 минут до 3 часов

Опишем реальный сценарий нарастания проблемы на протяжении 6 месяцев.

Январь: всё работает нормально

Таблица events: 500 GB, 4000 файлов по 125 MB
Spark запрос: SELECT COUNT(*) GROUP BY event_type
- Listing: 4 ListObjectsV2 запроса → 80 мс
- Reading: 4000 × 20 мс TTFB = 80 сек + данные
- Общее время: 2 минуты ← норма

Апрель: streaming без контроля файлов

# Streaming запущен в феврале с настройками:
df.writeStream \
    .trigger(processingTime="1 minute") \  # 60 batch/час
    .format("parquet") \
    .partitionBy("event_date", "event_hour") \  # 24 × 90 дней = 2160 партиций
    .start("s3a://bucket/events/")

# 60 batch/час × 24 часа × 90 дней = 129 600 batch
# Каждый batch: 50 активных воркеров → 50 файлов
# Итого: 129 600 × 50 = 6 480 000 новых мелких файлов за 3 месяца!
Апрель: 6 484 000 файлов (в основном 100 KB–1 MB)
Spark запрос: SELECT COUNT(*) GROUP BY event_type
- Listing: 6 484 ListObjectsV2 запросов → 130 СЕКУНД только на листинг
- OOM на Driver при 8GB памяти
- Запрос падает с java.lang.OutOfMemoryError
- После увеличения до 32 GB Driver: запрос выполняется 3 часа

Решение: применение комплекса мер

# Шаг 1: Немедленная миграция на Delta Lake (устраняет LIST при планировании):
spark.sql("""
    CONVERT TO DELTA parquet.`s3a://bucket/events/`
    PARTITIONED BY (event_date STRING, event_hour INT)
""")

# Шаг 2: OPTIMIZE для консолидации существующих файлов:
spark.sql("""
    OPTIMIZE delta.`s3a://bucket/events`
""")
# Время: 4 часа (однократно)
# Результат: 6.4M файлов → 2000 файлов по 256 MB

# Шаг 3: Перенастройка streaming:
df.writeStream \
    .trigger(processingTime="15 minutes") \  # вместо 1 минуты
    .foreachBatch(lambda df, id: (
        df.coalesce(4)
          .write.mode("append").format("delta")
          .partitionBy("event_date")
          .save("s3a://bucket/events/")
    )) \
    .start()

# Шаг 4: Ежедневная компакция через Airflow:
# OPTIMIZE delta.`s3a://bucket/events`
# WHERE event_date = yesterday
# ZORDER BY (user_id, event_type)
После оптимизации:
- Файлов: 2 000 (вместо 6 484 000)
- Spark запрос: 2 минуты 15 секунд (вместо 3 часов)
- API стоимость: $0.002/запрос (вместо $2.59/запрос)
- Ускорение: 80×
- Снижение API стоимости: 1 295×

Чеклист: предотвращение Small Files Problem

  1. Целевой размер файла - 128–512 MB. Файлы меньше 64 MB - подозрительно.

  2. Партиционировать по низкой кардинальности: день, месяц, регион (< 10 000 уникальных значений). Высокая кардинальность → bucket partitioning.

  3. Streaming trigger - не менее 5–15 минут. Меньше trigger = больше файлов. Использовать foreachBatch + coalesce.

  4. Запретить repartition(N) с большим N без контроля выходного размера. После repartition(1000) - coalesce до нужного числа файлов.

  5. Ежедневная компакция - OPTIMIZE (Delta) или rewrite_data_files (Iceberg) для партиций, в которые пишет streaming.

  6. Регулярный аудит - len(df.inputFiles()) перед деплоем нового пайплайна. Если > 10 000 на таблицу - нужна компакция.

  7. Lakehouse-форматы (Delta/Iceberg) - первый шаг при проблемах с LIST latency. Устраняют S3 LIST из критического пути планирования.

  8. AQE - включить spark.sql.adaptive.enabled=true для автоматического объединения мелких shuffle партиций.

  9. Driver память - при работе с большим числом файлов увеличить spark.driver.memory. Симптом OOM: java.lang.OutOfMemoryError при листинге.

  10. VACUUM / expire_snapshots - удалять старые файлы после компакции. Без этого диск S3 заполняется дважды (старые мелкие + новые крупные).