Small Files на S3: LIST latency и ограничения при большом числе файлов
Small Files Problem - главный инфраструктурный капкан при переходе на S3. Разбираем физику HTTP-накладных расходов, Scheduler Overhead Spark, стратегии компакции и как Lakehouse-форматы решают проблему метаданных.
Почему 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 запрос к удалённому веб-сервису. Каждый запрос включает:
- DNS-резолюцию (если соединение новое): ~5 мс.
- TCP-хэндшейк: ~1–2 RTT.
- TLS-хэндшейк: ~1–2 RTT дополнительно.
- Отправку HTTP-запроса.
- Ожидание первого байта ответа (TTFB): 5–50 мс.
- Получение данных по сети.
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 перед запуском задачи должен «открыть» все входные файлы и построить план:
- Рекурсивный листинг директорий (ListObjectsV2).
- Получение FileStatus для каждого файла (size, modification_time).
- Разбивка файлов на InputSplits (один Split = один таск по умолчанию).
- Сериализация 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¶
-
Целевой размер файла - 128–512 MB. Файлы меньше 64 MB - подозрительно.
-
Партиционировать по низкой кардинальности: день, месяц, регион (< 10 000 уникальных значений). Высокая кардинальность → bucket partitioning.
-
Streaming trigger - не менее 5–15 минут. Меньше trigger = больше файлов. Использовать
foreachBatch+coalesce. -
Запретить repartition(N) с большим N без контроля выходного размера. После repartition(1000) - coalesce до нужного числа файлов.
-
Ежедневная компакция - OPTIMIZE (Delta) или rewrite_data_files (Iceberg) для партиций, в которые пишет streaming.
-
Регулярный аудит -
len(df.inputFiles())перед деплоем нового пайплайна. Если > 10 000 на таблицу - нужна компакция. -
Lakehouse-форматы (Delta/Iceberg) - первый шаг при проблемах с LIST latency. Устраняют S3 LIST из критического пути планирования.
-
AQE - включить
spark.sql.adaptive.enabled=trueдля автоматического объединения мелких shuffle партиций. -
Driver память - при работе с большим числом файлов увеличить
spark.driver.memory. Симптом OOM:java.lang.OutOfMemoryErrorпри листинге. -
VACUUM / expire_snapshots - удалять старые файлы после компакции. Без этого диск S3 заполняется дважды (старые мелкие + новые крупные).