Spark

Junior

Проблема мелких файлов: почему ваша Spark-задача тратит больше времени на I/O, чем на вычисления

Задача может упираться в I/O, но выглядеть упирающейся в вычисления на любом дашборде, который отслеживает только CPU — потому что узкое место в открытии файлов, а не в чтении байтов.

Читать на:
Расходы на открытие файла Расходы на планирование Реальное время чтения
Куда уходит время в стейдже, затронутом проблемой мелких файлов, примерно — иллюстративно, не измерено.

Ниже определённого размера файла фиксированные накладные расходы на его открытие — поиск метаданных, планирование задачи, настройка кодека — обходятся дороже, чем чтение его реального содержимого.

Как получаются тысячи крошечных файлов

Обычные источники: партиционирование записи по колонке с высокой кардинальностью (каждое значение партиции получает свой маленький файл), стриминговые задачи, коммитящие мелкие микро-батчи с коротким интервалом триггера, или shuffle-стейдж с гораздо большим числом выходных партиций, чем реально нужно данным. Ни одно из этого не является ошибкой в строгом смысле — часто это правильное решение для задачи, которая пишет данные — но последующему чтению это оставляет беспорядок.

Почему это стоит дороже, чем подсказывает число байт

Каждый файл, даже крошечный, всё равно требует полноценной задачи: запланировать её, открыть файл (реальный сетевой round trip на S3 или GCS), прочитать footer или header, затем закрыть. При 10 000 крошечных файлов эти накладные расходы на файл могут доминировать над общим временем стейджа, хотя реальные данные едва заполняют память нескольких executor'ов — задача выглядит упирающейся в вычисления на графике CPU, потому что CPU занят планированием и учётом I/O, а не реальной работой.

# примерный средний размер файла для пути вывода avg_file_size = total_bytes_written / file_count # заметно меньше ~128МБ (Parquet) стоит расследовать

Фикс, который действительно решает проблему

coalesce() или repartition() прямо перед записью, с целью примерно 128МБ-1ГБ на выходной файл в зависимости от формата и движка, который будет читать данные дальше. Разовый фикс на стороне записи дешевле, чем каждая последующая задача, повторно платящая налог на открытие файлов — компактификация постфактум тоже работает, но это повторяющиеся расходы, которых правильно настроенная запись избегает полностью.

Посмотрите, как это выглядит на вашем собственном пайплайне.

Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.