Compaction стратегии: Spark job, Delta OPTIMIZE, Iceberg rewrite_data_files

Глубокий разбор стратегий компакции файлов в object storage: ручной Spark job, Delta Lake OPTIMIZE/VACUUM, Iceberg rewrite_data_files/expire_snapshots, автоматическая компакция в стриминге и планирование через Airflow.

storage optimization

Компакция как операционная практика

Предыдущий урок показал, почему маленькие файлы - это катастрофа для производительности: 10 000 файлов по 13 КБ читаются в 155 раз медленнее, чем один файл 128 МБ того же объёма. Причина - Time To First Byte (TTFB) в 130 мс на каждый HTTP GET-запрос, который не зависит от размера файла.

Компакция (compaction) - это процесс объединения множества мелких файлов в меньшее количество оптимальных по размеру файлов. По сути, это операция «прочитать N маленьких файлов, записать M больших файлов» (где N >> M), выполняемая периодически или по триггеру.

Компакция - не разовая операция, а непрерывная операционная практика. Данные неизбежно фрагментируются: стриминг пишет файл каждые 30 секунд, ETL делает append без coalesce, retry-механизмы создают дублирующие файлы. Без регулярной компакции любая система хранения данных деградирует.

Цели компакции:

  1. Сокращение числа файлов - уменьшение числа HTTP GET-запросов при чтении
  2. Оптимизация размера файлов - попадание в диапазон 128–512 МБ (sweet spot для Parquet)
  3. Кластеризация данных - во время компакции можно отсортировать данные по часто используемым в фильтрах полям, кардинально ускорив будущие запросы
  4. Уменьшение расходов на S3 API - меньше файлов = меньше LIST и GET запросов = меньше счёт

Математика оптимального размера файла

Прежде чем переходить к инструментам компакции, важно понять, какой размер файла является целевым. Слишком мелкие файлы - проблема TTFB. Слишком крупные файлы - своя проблема.

Почему файл не должен быть больше 512–1024 МБ:

  • Один Spark executor читает один Row Group за раз. Row Group по умолчанию - 128 МБ. Файл 1 ГБ содержит ~8 Row Groups, и если Spark'у нужны данные только из одного Row Group, он всё равно скачивает все 8 (если нет хорошей статистики и предиктных фильтров).
  • Крупные файлы хуже параллелизируются. Spark создаёт один input split на каждые spark.sql.files.maxPartitionBytes (по умолчанию 128 МБ). Файл 10 ГБ → 80 splits → 80 tasks, но они не будут работать параллельно, если их больше, чем executor'ов.
  • При повреждении файла теряется больше данных.

Практические диапазоны:

Тип нагрузки Целевой размер файла Row Group size Обоснование
OLAP запросы (аналитика) 256–512 МБ 128–256 МБ Больше Row Groups = лучше pruning
Смешанная нагрузка 128–256 МБ 128 МБ Баланс параллелизма и TTFB
Стриминг (частые writes) 64–128 МБ 64 МБ Частые compaction, меньше overhead
Исторические архивы 512 МБ–1 ГБ 256 МБ Редко читаются, максимальная компрессия

Формула целевого числа файлов после компакции:

import math

def calculate_target_file_count(
    total_size_bytes: int,
    target_file_size_bytes: int = 256 * 1024 * 1024,  # 256 МБ
    min_files: int = 1
) -> int:
    """
    Рассчитывает оптимальное число файлов после компакции.

    total_size_bytes: суммарный размер всех файлов до компакции
    target_file_size_bytes: целевой размер каждого выходного файла
    min_files: минимальное число файлов (не схлопывать в 1)
    """
    if total_size_bytes <= 0:
        return min_files

    raw_count = math.ceil(total_size_bytes / target_file_size_bytes)
    return max(raw_count, min_files)


# Пример расчёта
total_gb = 45  # 45 ГБ данных в разделе
total_bytes = total_gb * 1024**3

target_files = calculate_target_file_count(
    total_size_bytes=total_bytes,
    target_file_size_bytes=256 * 1024 * 1024  # 256 МБ
)
print(f"45 ГБ данных → {target_files} файлов по 256 МБ")
# 45 ГБ данных → 180 файлов по 256 МБ

Ручная компакция через Spark Job

Самый базовый способ компакции - написать Spark job вручную. Это даёт полный контроль над процессом, но имеет ряд ограничений.

Простейшая компакция: repartition + write

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

# Простейшая компакция: читаем → переразбиваем → перезаписываем
(
    spark.read.parquet("s3a://bucket/data/events/")
    .repartition(180)          # 45 ГБ / 256 МБ = 180 файлов
    .write
    .mode("overwrite")
    .parquet("s3a://bucket/data/events/")
)

На первый взгляд просто. Но у этого кода критическая проблема: читаем из того же пути, куда пишем. Когда Spark начинает запись в режиме overwrite, он сначала удаляет существующие данные, а потом пишет новые. Если что-то пойдёт не так посередине - данные потеряны.

Правильный подход - Write-to-Temp, then Swap:

import math
from dataclasses import dataclass
from pyspark.sql import SparkSession, DataFrame
from typing import Optional


@dataclass
class CompactionConfig:
    """Конфигурация задачи компакции."""
    source_path: str
    output_path: str              # может совпадать с source_path при использовании tmp
    temp_path: str                # временный путь для записи
    target_file_size_mb: int = 256
    sort_columns: Optional[list] = None   # сортировать данные во время компакции?
    partition_columns: Optional[list] = None  # колонки партиционирования
    compression: str = "snappy"


def get_table_size_bytes(spark: SparkSession, path: str) -> int:
    """
    Получает суммарный размер всех файлов по пути.
    Использует Hadoop FileSystem API через Spark.
    """
    sc = spark.sparkContext
    hadoop_conf = sc._jvm.org.apache.hadoop.conf.Configuration()
    fs = sc._jvm.org.apache.hadoop.fs.FileSystem.get(
        sc._jvm.java.net.URI.create(path), hadoop_conf
    )
    content_summary = fs.getContentSummary(
        sc._jvm.org.apache.hadoop.fs.Path(path)
    )
    return content_summary.getLength()


def calculate_target_partitions(
    total_bytes: int,
    target_file_size_mb: int = 256,
    min_partitions: int = 1
) -> int:
    """Рассчитывает оптимальное число партиций Spark (= число файлов)."""
    target_bytes = target_file_size_mb * 1024 * 1024
    raw = math.ceil(total_bytes / target_bytes)
    return max(raw, min_partitions)


def compact_dataframe(
    df: DataFrame,
    config: CompactionConfig,
    total_bytes: int
) -> DataFrame:
    """
    Применяет компакцию к DataFrame:
    1. Рассчитывает оптимальное число партиций
    2. Опционально сортирует для кластеризации
    3. Возвращает готовый к записи DataFrame
    """
    target_partitions = calculate_target_partitions(
        total_bytes=total_bytes,
        target_file_size_mb=config.target_file_size_mb
    )
    print(f"Компакция: {total_bytes / 1024**3:.2f} ГБ → {target_partitions} файлов")

    if config.sort_columns:
        # sortWithinPartitions: сортировка внутри каждой партиции
        # НЕ то же самое, что orderBy (не вызывает глобальный shuffle)
        df = df.sortWithinPartitions(*config.sort_columns)

    return df.repartition(target_partitions)


def run_manual_compaction(spark: SparkSession, config: CompactionConfig) -> None:
    """
    Безопасная компакция через временный путь:
    1. Читаем источник
    2. Пишем во временный путь
    3. Атомарно переименовываем (или Hadoop rename на HDFS)
    """
    sc = spark.sparkContext
    hadoop_conf = sc._jvm.org.apache.hadoop.conf.Configuration()
    fs = sc._jvm.org.apache.hadoop.fs.FileSystem.get(
        sc._jvm.java.net.URI.create(config.source_path), hadoop_conf
    )

    # Шаг 1: измеряем объём данных
    total_bytes = get_table_size_bytes(spark, config.source_path)
    print(f"Источник: {config.source_path}")
    print(f"Объём: {total_bytes / 1024**3:.2f} ГБ")

    # Шаг 2: читаем данные
    df = spark.read.parquet(config.source_path)

    # Шаг 3: применяем компакцию
    compacted_df = compact_dataframe(df, config, total_bytes)

    # Шаг 4: пишем во временный путь
    writer = (
        compacted_df
        .write
        .mode("overwrite")
        .option("compression", config.compression)
    )

    if config.partition_columns:
        writer = writer.partitionBy(*config.partition_columns)

    writer.parquet(config.temp_path)
    print(f"Записано во временный путь: {config.temp_path}")

    # Шаг 5: удаляем старые данные и переименовываем
    # ВНИМАНИЕ: на S3 rename = copy + delete, это не атомарная операция!
    # Для production используйте Delta Lake или Iceberg
    src_path = sc._jvm.org.apache.hadoop.fs.Path(config.temp_path)
    dst_path = sc._jvm.org.apache.hadoop.fs.Path(config.output_path)

    if fs.exists(dst_path):
        fs.delete(dst_path, True)  # рекурсивное удаление

    fs.rename(src_path, dst_path)
    print(f"Компакция завершена: {config.output_path}")

Компакция партиционированных таблиц

Когда таблица партиционирована (например, по event_date), нужно компактировать каждый раздел отдельно. Иначе Spark перемешает данные из разных партиций в одни файлы, и партиционирование потеряет смысл.

from datetime import date, timedelta
from pyspark.sql import SparkSession


def compact_partitioned_table(
    spark: SparkSession,
    base_path: str,
    partition_column: str,
    partitions: list,
    target_file_size_mb: int = 256,
    sort_columns: Optional[list] = None
) -> None:
    """
    Компактирует партиционированную таблицу, обрабатывая каждый раздел отдельно.

    Пример вызова:
        compact_partitioned_table(
            spark=spark,
            base_path="s3a://bucket/events/",
            partition_column="event_date",
            partitions=["2024-01-01", "2024-01-02"],
            sort_columns=["user_id", "event_time"]
        )
    """
    for partition_value in partitions:
        partition_path = f"{base_path}/{partition_column}={partition_value}"
        temp_path = f"{base_path}/_compaction_tmp/{partition_column}={partition_value}"

        print(f"\nКомпактируем: {partition_column}={partition_value}")

        config = CompactionConfig(
            source_path=partition_path,
            output_path=partition_path,
            temp_path=temp_path,
            target_file_size_mb=target_file_size_mb,
            sort_columns=sort_columns
        )

        run_manual_compaction(spark, config)

    # Очищаем временную директорию
    sc = spark.sparkContext
    hadoop_conf = sc._jvm.org.apache.hadoop.conf.Configuration()
    fs = sc._jvm.org.apache.hadoop.fs.FileSystem.get(
        sc._jvm.java.net.URI.create(base_path), hadoop_conf
    )
    tmp_path = sc._jvm.org.apache.hadoop.fs.Path(f"{base_path}/_compaction_tmp")
    if fs.exists(tmp_path):
        fs.delete(tmp_path, True)
    print("\nКомпакция всех партиций завершена")


# Компактируем последние 7 дней
today = date.today()
recent_partitions = [
    str(today - timedelta(days=i))
    for i in range(7)
]

compact_partitioned_table(
    spark=spark,
    base_path="s3a://bucket/events/",
    partition_column="event_date",
    partitions=recent_partitions,
    target_file_size_mb=256,
    sort_columns=["user_id", "event_type"]
)

Ограничения ручной компакции

Ручная компакция - наивный подход, у которого есть серьёзные проблемы:

  1. Отсутствие атомарности на S3. Hadoop rename на S3 = copy + delete. Если процесс упадёт после копирования, но до удаления - получим дубли. Если после удаления, но до завершения копирования - потеряем данные.

  2. Нет защиты от concurrent читателей. Пока идёт компакция, читатели видят смесь старых и новых файлов.

  3. Нет отката. Если компакция записала данные неверно, отката нет - нужно перечитывать из резервной копии.

  4. Нет метаданных о размерах файлов. Нужно каждый раз сканировать весь путь через getContentSummary, что само по себе нагружает S3 API.

Именно эти проблемы решают Delta Lake и Apache Iceberg, предоставляя встроенные механизмы компакции.


Delta Lake: OPTIMIZE - умная компакция с транзакциями

Delta Lake хранит все метаданные о файлах в _delta_log/. Каждая операция (добавление, удаление файлов) записывается как JSON-файл транзакции. Это делает возможной атомарную компакцию.

Этот механизм гарантирует: если OPTIMIZE упадёт в любой момент до записи финальной транзакции в _delta_log/, старые файлы остаются активными. Читатели ничего не заметят.

Команда OPTIMIZE: синтаксис

-- Простой OPTIMIZE всей таблицы
OPTIMIZE events;

-- OPTIMIZE с фильтром по партиции (рекомендуется для больших таблиц)
-- Compact только один раздел - быстрее и экономичнее
OPTIMIZE events
WHERE event_date = '2024-01-15';

-- OPTIMIZE с диапазоном дат
OPTIMIZE events
WHERE event_date BETWEEN '2024-01-01' AND '2024-01-31';

-- OPTIMIZE с Z-Order - компакция + кластеризация данных
-- Z-Order физически располагает похожие значения рядом в файле
-- Это позволяет Row Group Filtering пропускать больше данных
OPTIMIZE events
ZORDER BY (user_id, event_type);

-- OPTIMIZE конкретной партиции + Z-Order
OPTIMIZE events
WHERE event_date = '2024-01-15'
ZORDER BY (user_id, event_type);

OPTIMIZE через Python API

from delta import DeltaTable
from pyspark.sql import SparkSession


def optimize_delta_table(
    spark: SparkSession,
    table_path: str,
    partition_filter: Optional[str] = None,
    zorder_columns: Optional[list] = None
) -> dict:
    """
    Выполняет OPTIMIZE для Delta Lake таблицы.

    Возвращает словарь с метриками:
    - numFilesAdded: число добавленных файлов (после компакции)
    - numFilesRemoved: число удалённых файлов (до компакции)
    - filesAdded.avg: средний размер новых файлов в байтах
    - filesRemoved.avg: средний размер старых файлов в байтах
    """
    dt = DeltaTable.forPath(spark, table_path)

    if partition_filter and zorder_columns:
        result = (
            dt.optimize()
            .where(partition_filter)
            .executeZOrderBy(*zorder_columns)
        )
    elif partition_filter:
        result = (
            dt.optimize()
            .where(partition_filter)
            .executeCompaction()
        )
    elif zorder_columns:
        result = (
            dt.optimize()
            .executeZOrderBy(*zorder_columns)
        )
    else:
        result = dt.optimize().executeCompaction()

    # result.show() - выводит метрики компакции
    metrics = result.first().asDict()
    return metrics


# Использование
metrics = optimize_delta_table(
    spark=spark,
    table_path="s3a://bucket/events/",
    partition_filter="event_date = '2024-01-15'",
    zorder_columns=["user_id", "event_type"]
)

print(f"Файлов до: {metrics['numFilesRemoved']}")
print(f"Файлов после: {metrics['numFilesAdded']}")
if metrics.get('filesRemoved') and metrics['filesRemoved'].get('avg'):
    print(f"Средний размер ДО: {metrics['filesRemoved']['avg'] / 1024**2:.1f} МБ")
if metrics.get('filesAdded') and metrics['filesAdded'].get('avg'):
    print(f"Средний размер ПОСЛЕ: {metrics['filesAdded']['avg'] / 1024**2:.1f} МБ")

Конфигурация целевого размера файла в Delta

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog",
            "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    # Целевой размер файла после OPTIMIZE (default: 1 ГБ в Delta 2.x)
    .config("spark.databricks.delta.optimize.maxFileSize",
            str(256 * 1024 * 1024)) \
    # Минимальный размер файла для участия в компакции
    # Файлы крупнее этого порога пропускаются
    .config("spark.databricks.delta.optimize.minFileSize",
            str(32 * 1024 * 1024)) \
    .getOrCreate()

# Или через ALTER TABLE (per-table конфигурация)
spark.sql("""
    ALTER TABLE events
    SET TBLPROPERTIES (
        'delta.targetFileSize' = '268435456',     -- 256 МБ
        'delta.tuneFileSizesForRewrites' = 'true'
    )
""")

Delta Auto-Optimize: оптимизированные записи и авто-компакция

Delta Lake поддерживает автоматическую оптимизацию при записи:

# Optimized Writes: Spark сам коалесцирует партиции перед записью
# Вместо 200 мелких файлов пишет 10 больших
# Работает для append и overwrite
spark.conf.set("spark.databricks.delta.optimizeWrite.enabled", "true")

# Auto Compaction: после каждой записи Delta проверяет,
# не накопилось ли слишком много мелких файлов, и компактирует автоматически
# Триггерится асинхронно в фоне
spark.conf.set("spark.databricks.delta.autoCompact.enabled", "true")

# Задаём порог авто-компакции: если в партиции больше N файлов,
# запускаем компакцию
spark.conf.set("spark.databricks.delta.autoCompact.minNumFiles", "50")

# Пишем данные - Delta автоматически оптимизирует запись и запускает компакцию
df.write \
    .format("delta") \
    .mode("append") \
    .option("mergeSchema", "true") \
    .save("s3a://bucket/events/")

Когда использовать Auto-Optimize:

  • Стриминговые пайплайны, где ручная компакция сложна в планировании
  • Таблицы с хаотичными паттернами записи (разные команды, разные размеры батчей)
  • Dev/staging окружения, где операционный overhead нежелателен

Когда лучше ручной OPTIMIZE:

  • Production таблицы, где нужна предсказуемость (OPTIMIZE в заданное время)
  • Когда нужен Z-Order - Auto-Compact не делает Z-Ordering
  • Когда хочется контролировать нагрузку (OPTIMIZE ночью, не во время пиковых запросов)

Delta Lake: VACUUM - физическое удаление устаревших файлов

После OPTIMIZE старые мелкие файлы физически остаются на S3. Delta лишь помечает их как RemoveFile в Delta Log, но сами объекты никуда не исчезают. Это сделано намеренно: Delta поддерживает time travel (запросы к историческим версиям таблицы) и concurrent readers (транзакция, начавшаяся до OPTIMIZE, должна видеть старые файлы).

Синтаксис VACUUM

-- Удалить все файлы, которые Delta не отслеживает и которым > 7 дней (168 часов)
VACUUM events RETAIN 168 HOURS;

-- Режим DRY RUN: показывает, что БЫЛО БЫ удалено, но не удаляет
-- Рекомендуется запускать перед первым VACUUM на production
VACUUM events RETAIN 168 HOURS DRY RUN;

-- ОПАСНО: удалить файлы старше 0 часов
-- Ломает time travel и может сломать concurrent транзакции
-- Delta блокирует это по умолчанию (retentionDurationCheck)
SET spark.databricks.delta.retentionDurationCheck.enabled = false;
VACUUM events RETAIN 0 HOURS;  -- крайне не рекомендуется в production!

Почему нельзя ставить RETAIN меньше 7 дней

168 часов (7 дней) - это минимум, рекомендованный Delta. Причины:

  1. Concurrent readers: транзакция, начавшаяся 6 часов назад, читает файлы по состоянию на момент своего начала. Если VACUUM удалит эти файлы, транзакция упадёт с FileNotFoundException.

  2. Time travel: SELECT * FROM events TIMESTAMP AS OF '2024-01-14' требует файлов, которые существовали в тот момент. Без них - ошибка.

  3. Backup window: если сегодня обнаружили баг в ETL, который портит данные начиная с позавчера - 7-дневное окно позволяет откатиться к чистым данным.

  4. Streaming checkpoints: стриминговые задания используют offset'ы, указывающие на конкретные версии Delta. Уменьшение retention может сломать эти checkpoint'ы.

from delta import DeltaTable

def vacuum_delta_table(
    spark: SparkSession,
    table_path: str,
    retain_hours: int = 168,   # 7 дней
    dry_run: bool = True       # по умолчанию только показываем
) -> None:
    """Безопасное удаление устаревших файлов через Delta VACUUM."""
    dt = DeltaTable.forPath(spark, table_path)

    if dry_run:
        print("=== DRY RUN: файлы, которые будут удалены ===")
        dt.vacuum(retentionHours=retain_hours)
        # В режиме dry_run нужен spark.conf.set ниже
    else:
        if retain_hours < 168:
            print(f"ПРЕДУПРЕЖДЕНИЕ: retain_hours={retain_hours} < 168 (7 дней)!")
            print("Это может сломать concurrent транзакции и time travel")
            # В production НИКОГДА не уменьшать ниже 168
            return

        dt.vacuum(retentionHours=retain_hours)
        print(f"VACUUM завершён. Удалены файлы старше {retain_hours}ч ({retain_hours//24} дней)")


# Правильный production workflow:
# 1. Сначала OPTIMIZE
optimize_delta_table(
    spark, "s3a://bucket/events/",
    partition_filter="event_date = '2024-01-15'"
)
# 2. Потом VACUUM (через 7+ дней)
vacuum_delta_table(
    spark, "s3a://bucket/events/",
    retain_hours=168,
    dry_run=False
)

Apache Iceberg: rewrite_data_files

Iceberg предоставляет аналогичный механизм через SQL-процедуру rewrite_data_files. В отличие от Delta, Iceberg работает с концепцией snapshot'ов: каждая операция создаёт новый snapshot, указывающий на новый набор файлов. Старый snapshot и его файлы остаются доступными до явного expire_snapshots.

Атомарность гарантирована механизмом snapshot'ов: если rewrite_data_files упадёт до создания Snapshot 1, Snapshot 0 остаётся текущим. Читатели ничего не заметят.

Синтаксис rewrite_data_files

-- Базовая компакция с binpack-стратегией
-- binpack: упаковывает мелкие файлы в крупные без изменения порядка строк
CALL catalog.system.rewrite_data_files(
    table => 'mydb.events'
);

-- Компакция только конкретной партиции
CALL catalog.system.rewrite_data_files(
    table => 'mydb.events',
    where => 'event_date = ''2024-01-15'''
);

-- Компакция с настройкой целевого размера файла
CALL catalog.system.rewrite_data_files(
    table => 'mydb.events',
    options => map(
        'target-file-size-bytes', '268435456'   -- 256 МБ
    )
);

-- Компакция с фильтром по размеру: компактировать только файлы < 128 МБ
-- Крупные файлы (уже оптимального размера) пропускаются
CALL catalog.system.rewrite_data_files(
    table => 'mydb.events',
    options => map(
        'target-file-size-bytes',   '268435456',  -- цель: 256 МБ
        'min-file-size-bytes',      '1048576',    -- компактировать файлы > 1 МБ...
        'max-file-size-bytes',      '134217728'   -- ...и < 128 МБ (мелкие)
    )
);

-- Sort-стратегия: перетасовывает данные так, чтобы похожие значения
-- оказались в одних Row Groups - ускоряет запросы с фильтрами
CALL catalog.system.rewrite_data_files(
    table => 'mydb.events',
    strategy => 'sort',
    sort_order => 'user_id ASC NULLS LAST, event_date ASC NULLS LAST'
);

-- Partial progress: для очень больших таблиц
-- Делает промежуточные коммиты каждые N групп файлов
-- Если компакция упадёт на 50% - половина работы уже закоммичена
CALL catalog.system.rewrite_data_files(
    table => 'mydb.events',
    options => map(
        'target-file-size-bytes',       '268435456',
        'partial_progress.enabled',     'true',
        'partial_progress.max_commits', '10'
    )
);

rewrite_data_files через Python API

from pyspark.sql import SparkSession


def iceberg_compact(
    spark: SparkSession,
    catalog: str,
    database: str,
    table: str,
    partition_filter: Optional[str] = None,
    strategy: str = "binpack",
    sort_order: Optional[str] = None,
    target_file_size_mb: int = 256,
    min_file_size_mb: int = 1,
    max_file_size_mb: int = 128,
    partial_progress: bool = False,
    max_commits: int = 10
) -> None:
    """
    Выполняет компакцию Iceberg-таблицы через rewrite_data_files.

    strategy:
        'binpack' - упаковывает мелкие файлы без сортировки (быстро)
        'sort'    - сортирует данные (медленнее, но лучше для запросов)
    """
    options = {
        "target-file-size-bytes": str(target_file_size_mb * 1024 * 1024),
        "min-file-size-bytes": str(min_file_size_mb * 1024 * 1024),
        "max-file-size-bytes": str(max_file_size_mb * 1024 * 1024),
    }

    if partial_progress:
        options["partial_progress.enabled"] = "true"
        options["partial_progress.max_commits"] = str(max_commits)

    # Строим SQL опции как Spark map
    options_sql = ", ".join(
        f"'{k}', '{v}'" for k, v in options.items()
    )

    if strategy == "sort" and sort_order:
        sql = f"""
            CALL {catalog}.system.rewrite_data_files(
                table => '{database}.{table}',
                strategy => 'sort',
                sort_order => '{sort_order}',
                options => map({options_sql})
                {f", where => '{partition_filter}'" if partition_filter else ""}
            )
        """
    else:
        sql = f"""
            CALL {catalog}.system.rewrite_data_files(
                table => '{database}.{table}',
                strategy => '{strategy}',
                options => map({options_sql})
                {f", where => '{partition_filter}'" if partition_filter else ""}
            )
        """

    result = spark.sql(sql)
    row = result.first()

    if row:
        print(f"rewrite_data_files результат:")
        print(f"  Файлов переписано: {row['rewritten_data_files_count']}")
        print(f"  Файлов добавлено:  {row['added_data_files_count']}")
        print(f"  Байт переписано:   {row['rewritten_bytes_count'] / 1024**3:.2f} ГБ")


# Binpack - ежедневная компакция вчерашнего раздела
iceberg_compact(
    spark=spark,
    catalog="iceberg",
    database="analytics",
    table="events",
    partition_filter="event_date = '2024-01-14'",
    strategy="binpack",
    target_file_size_mb=256
)

# Sort - еженедельная компакция с сортировкой для ускорения запросов
iceberg_compact(
    spark=spark,
    catalog="iceberg",
    database="analytics",
    table="events",
    strategy="sort",
    sort_order="user_id ASC NULLS LAST, event_date ASC NULLS LAST",
    target_file_size_mb=256,
    partial_progress=True
)

Apache Iceberg: expire_snapshots и remove_orphan_files

После нескольких rewrite_data_files у таблицы накапливается множество snapshot'ов. Каждый snapshot ссылается на данные, которые были актуальны на момент его создания. Пока snapshot жив - файлы, которые он видит, нельзя удалять (они ещё нужны).

expire_snapshots удаляет старые snapshot'ы и освобождает файлы, на которые они ссылались, для последующей физической очистки.

-- Удалить snapshot'ы старше указанного времени
CALL catalog.system.expire_snapshots(
    table => 'mydb.events',
    older_than => TIMESTAMP '2024-01-08 00:00:00'
);

-- Оставить минимум N последних snapshot'ов (независимо от возраста)
CALL catalog.system.expire_snapshots(
    table => 'mydb.events',
    older_than => TIMESTAMP '2024-01-08 00:00:00',
    retain_last => 5
);

-- Удалить snapshot'ы И сразу же физически удалить файлы
-- (иначе нужен отдельный вызов remove_orphan_files)
CALL catalog.system.expire_snapshots(
    table => 'mydb.events',
    older_than => TIMESTAMP '2024-01-08 00:00:00',
    retain_last => 5
    -- delete_data = true  -- осторожно! удаляет файлы немедленно
);

-- Удалить файлы-сироты (orphan files)
-- Это файлы, которые существуют в S3, но не упоминаются ни в одном snapshot'е
-- Может возникнуть после упавшей записи, ручных операций с файлами, etc.
CALL catalog.system.remove_orphan_files(
    table => 'mydb.events',
    older_than => TIMESTAMP '2024-01-08 00:00:00'   -- файлы старше 7 дней
);
from datetime import datetime, timedelta
from pyspark.sql import SparkSession


def iceberg_maintenance(
    spark: SparkSession,
    catalog: str,
    database: str,
    table: str,
    retain_days: int = 7,
    retain_last_snapshots: int = 5
) -> None:
    """
    Полный maintenance цикл для Iceberg таблицы:
    1. expire_snapshots (удаляем старые snapshot'ы)
    2. remove_orphan_files (удаляем файлы-сироты)
    """
    cutoff = datetime.utcnow() - timedelta(days=retain_days)
    cutoff_str = cutoff.strftime("%Y-%m-%d %H:%M:%S")
    full_table = f"{database}.{table}"

    print(f"Iceberg Maintenance: {full_table}")
    print(f"Порог: {cutoff_str} (retain={retain_days} дней)")

    # Шаг 1: expire_snapshots
    result = spark.sql(f"""
        CALL {catalog}.system.expire_snapshots(
            table => '{full_table}',
            older_than => TIMESTAMP '{cutoff_str}',
            retain_last => {retain_last_snapshots}
        )
    """)
    row = result.first()
    if row:
        print(f"expire_snapshots: {row['deleted_data_files_count']} файлов,"
              f" {row['deleted_manifest_files_count']} манифестов удалено")

    # Шаг 2: remove_orphan_files
    # older_than должен быть >= cutoff у expire_snapshots
    orphan_result = spark.sql(f"""
        CALL {catalog}.system.remove_orphan_files(
            table => '{full_table}',
            older_than => TIMESTAMP '{cutoff_str}'
        )
    """)
    orphan_row = orphan_result.first()
    if orphan_row:
        print(f"remove_orphan_files: {orphan_row['orphan_file_location']} файлов удалено")

    print(f"Maintenance завершён: {full_table}")


# Запуск ежедневного maintenance (после ночной компакции)
iceberg_maintenance(
    spark=spark,
    catalog="iceberg",
    database="analytics",
    table="events",
    retain_days=7,
    retain_last_snapshots=10
)

Стратегии планирования компакции

Компакция - затратная операция: она читает все данные, перемешивает их и записывает заново. Важно правильно выбрать время и частоту запусков.

Стратегия 1: Расписание (Scheduled Compaction)

Самый простой подход: запускать компакцию раз в день (обычно ночью, в период минимальной нагрузки).

# Пример Airflow DAG для ночной компакции Delta Lake
from airflow import DAG
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta


with DAG(
    dag_id="delta_nightly_compaction",
    start_date=datetime(2024, 1, 1),
    schedule_interval="0 2 * * *",   # 02:00 UTC каждый день
    catchup=False,
    default_args={
        "owner": "data-platform",
        "retries": 1,
        "retry_delay": timedelta(minutes=10)
    }
) as dag:

    def get_yesterday_partitions(**context) -> list:
        """Возвращает партиции вчерашнего дня."""
        yesterday = context["ds"]   # YYYY-MM-DD вчерашний день
        return [yesterday]

    # Задача 1: получаем список партиций для компакции
    get_partitions = PythonOperator(
        task_id="get_partitions",
        python_callable=get_yesterday_partitions,
        provide_context=True
    )

    # Задача 2: компакция с OPTIMIZE
    optimize_events = SparkSubmitOperator(
        task_id="optimize_events",
        application="/jobs/compaction/optimize_events.py",
        conf={
            "spark.sql.extensions": "io.delta.sql.DeltaSparkSessionExtension",
            "spark.sql.catalog.spark_catalog":
                "org.apache.spark.sql.delta.catalog.DeltaCatalog",
        },
        application_args=[
            "--table-path", "s3a://bucket/events/",
            "--partition-date", "{{ ds }}"
        ]
    )

    # Задача 3: VACUUM (старше 7 дней)
    vacuum_events = SparkSubmitOperator(
        task_id="vacuum_events",
        application="/jobs/compaction/vacuum_events.py",
        application_args=[
            "--table-path", "s3a://bucket/events/",
            "--retain-hours", "168"
        ]
    )

    get_partitions >> optimize_events >> vacuum_events

Стратегия 2: Метрическая компакция (Metric-Based)

Более умный подход: запускать компакцию только тогда, когда метрики говорят «пора».

from pyspark.sql import SparkSession


def check_compaction_needed(
    spark: SparkSession,
    table_path: str,
    max_small_files: int = 1000,
    min_avg_size_mb: float = 32.0
) -> dict:
    """
    Проверяет, нужна ли компакция на основе метрик файлов.

    Возвращает словарь:
    {
        'needs_compaction': bool,
        'file_count': int,
        'avg_size_mb': float,
        'total_size_gb': float,
        'reason': str
    }
    """
    dt = DeltaTable.forPath(spark, table_path)
    detail = dt.detail().first()

    file_count = detail["numFiles"]
    total_bytes = detail["sizeInBytes"]
    avg_size_mb = (total_bytes / file_count / 1024**2) if file_count > 0 else 0

    reasons = []

    if file_count > max_small_files:
        reasons.append(f"Файлов: {file_count} > {max_small_files}")

    if avg_size_mb < min_avg_size_mb:
        reasons.append(f"Средний размер: {avg_size_mb:.1f} МБ < {min_avg_size_mb} МБ")

    needs_compaction = len(reasons) > 0

    return {
        "needs_compaction": needs_compaction,
        "file_count": file_count,
        "avg_size_mb": avg_size_mb,
        "total_size_gb": total_bytes / 1024**3,
        "reason": "; ".join(reasons) if reasons else "Компакция не нужна"
    }


def smart_compaction(spark: SparkSession, table_path: str) -> None:
    """Запускает компакцию только если нужно."""
    metrics = check_compaction_needed(spark, table_path)

    print(f"Таблица: {table_path}")
    print(f"Файлов: {metrics['file_count']}")
    print(f"Средний размер: {metrics['avg_size_mb']:.1f} МБ")
    print(f"Общий объём: {metrics['total_size_gb']:.2f} ГБ")

    if not metrics["needs_compaction"]:
        print(f"✓ Компакция не нужна: {metrics['reason']}")
        return

    print(f"✗ Нужна компакция: {metrics['reason']}")
    optimize_delta_table(spark, table_path)
    print("✓ Компакция завершена")

Стратегия 3: Компакция после стриминга (Streaming-Triggered)

Для стриминговых пайплайнов оптимально запускать компакцию после завершения каждого N-го микро-батча или по таймеру.

from pyspark.sql import SparkSession, DataFrame
from pyspark.sql.streaming import StreamingQuery


class StreamingCompactionManager:
    """
    Управляет компакцией для стриминговой таблицы.
    Запускает компакцию каждые N батчей или по таймеру.
    """

    def __init__(
        self,
        spark: SparkSession,
        output_path: str,
        compaction_interval_batches: int = 20,
        target_file_size_mb: int = 256
    ):
        self.spark = spark
        self.output_path = output_path
        self.compaction_interval = compaction_interval_batches
        self.target_file_size_mb = target_file_size_mb
        self.batch_count = 0

    def process_batch(self, df: DataFrame, batch_id: int) -> None:
        """
        foreachBatch handler:
        1. Пишет батч с coalesce (предотвращаем мелкие файлы)
        2. Каждые N батчей запускает OPTIMIZE
        """
        # Считаем размер батча для оптимального числа файлов
        row_count = df.count()
        # Предполагаем ~1 КБ на строку (подстраивайте под реальные данные)
        estimated_bytes = row_count * 1024
        target_files = calculate_target_partitions(
            total_bytes=estimated_bytes,
            target_file_size_mb=64   # для стриминга - файлы меньше
        )

        # Запись с coalesce: предотвращаем сотни мелких файлов
        (
            df
            .coalesce(target_files)
            .write
            .format("delta")
            .mode("append")
            .save(self.output_path)
        )

        self.batch_count += 1

        # Каждые N батчей - OPTIMIZE
        if self.batch_count % self.compaction_interval == 0:
            print(f"[Batch {batch_id}] Запуск компакции (каждые {self.compaction_interval} батчей)")
            optimize_delta_table(
                self.spark,
                self.output_path,
                # Компактируем только последние 2 часа данных
                partition_filter="event_date >= current_date() - interval 1 day"
            )


def start_streaming_with_compaction(spark: SparkSession) -> StreamingQuery:
    """Запускает стриминг с периодической компакцией."""
    manager = StreamingCompactionManager(
        spark=spark,
        output_path="s3a://bucket/events/",
        compaction_interval_batches=20,   # компактировать каждые 20 батчей
        target_file_size_mb=256
    )

    return (
        spark.readStream
        .format("kafka")
        .option("kafka.bootstrap.servers", "kafka:9092")
        .option("subscribe", "events")
        .load()
        .writeStream
        .foreachBatch(manager.process_batch)
        .trigger(processingTime="1 minute")
        .option("checkpointLocation", "s3a://bucket/checkpoints/events/")
        .start()
    )

Компакция с Z-Order кластеризацией

Компакция - это не только про размер файлов. Если во время компакции отсортировать данные по часто используемым в фильтрах колонкам, будущие запросы смогут пропускать целые Row Group'ы.

После сортировки по user_id Parquet-файл содержит Row Groups с чёткими диапазонами (min=1, max=500). Запрос WHERE user_id = 750 читает только один Row Group вместо всех - Row Group Filtering (predicate pushdown) работает эффективно.

Z-Order - это многомерная пространственная сортировка. В отличие от обычной однолонковой сортировки, Z-Order одновременно кластеризует по нескольким полям. Например, ZORDER BY (user_id, event_date) создаёт файлы, где данные по пользователю 123 за январь физически рядом с данными по пользователю 124 за январь - что невозможно при обычной сортировке.

# Delta Lake: OPTIMIZE с Z-Order
def optimize_with_zorder(
    spark: SparkSession,
    table_path: str,
    partition_filter: str,
    zorder_columns: list
) -> None:
    """
    Компакция с Z-Order кластеризацией.
    Z-Order оптимален для запросов с фильтрами по нескольким полям.
    """
    dt = DeltaTable.forPath(spark, table_path)

    print(f"OPTIMIZE ZORDER BY {zorder_columns}")
    result = (
        dt.optimize()
        .where(partition_filter)
        .executeZOrderBy(*zorder_columns)
    )

    row = result.first()
    if row:
        print(f"Завершено: {row['numFilesRemoved']}{row['numFilesAdded']} файлов")


# Пример: компакция таблицы events с Z-Order по user_id + event_type
optimize_with_zorder(
    spark=spark,
    table_path="s3a://bucket/events/",
    partition_filter="event_date >= '2024-01-01'",
    zorder_columns=["user_id", "event_type"]
)
-- Iceberg: sort strategy = аналог Z-Order
-- Сортирует данные во время rewrite_data_files
CALL catalog.system.rewrite_data_files(
    table => 'analytics.events',
    strategy => 'sort',
    sort_order => 'user_id ASC NULLS LAST, event_date ASC NULLS LAST'
);

-- Iceberg также поддерживает hidden partitioning с z-order:
ALTER TABLE analytics.events
WRITE ORDERED BY user_id, event_date;
-- Все будущие записи будут автоматически сортироваться

Сравнительная таблица подходов к компакции

Параметр Ручной Spark Job Delta OPTIMIZE Iceberg rewrite_data_files
Атомарность ❌ Нет (S3 rename) ✅ Delta Log транзакция ✅ Snapshot транзакция
Concurrent читатели ❌ Видят мешанину ✅ MVCC изоляция ✅ MVCC изоляция
Откат ❌ Нет ✅ Через time travel ✅ Через snapshot
Z-Order / Sort ➕ Вручную через sortWithinPartitions ✅ ZORDER BY ✅ sort strategy
Метрики компакции ❌ Нет ✅ Детальный отчёт ✅ Детальный отчёт
Partial progress ❌ Нет ➕ checkpoint-based ✅ partial_progress
Auto-Compaction ❌ Нет ✅ autoCompact ➕ Через процедуры
Фильтрация файлов ❌ Читает всё ✅ Только мелкие файлы ✅ min/max-file-size
Физическое удаление Немедленно VACUUM (7+ дней) expire_snapshots
Формат Любой Parquet Только Delta Только Iceberg

Практические Anti-Patterns

Anti-Pattern 1: OPTIMIZE всей таблицы каждый час

# НЕПРАВИЛЬНО: компактируем всю таблицу каждый час
# Для таблицы 10 ТБ: каждый час читаем и перезаписываем 10 ТБ данных
# Стоимость: 10 ТБ / час × 24 ч = 240 ТБ I/O в день
# Это парализует кластер и съест весь бюджет S3 API

schedule_interval = "0 * * * *"  # каждый час

def bad_hourly_optimize():
    optimize_delta_table(spark, "s3a://bucket/huge_events/")  # всё сразу!
# ПРАВИЛЬНО: компактируем только вчерашний раздел раз в день
schedule_interval = "0 2 * * *"  # 02:00 UTC каждый день

def good_daily_optimize():
    yesterday = (datetime.utcnow() - timedelta(days=1)).strftime("%Y-%m-%d")
    optimize_delta_table(
        spark,
        "s3a://bucket/events/",
        partition_filter=f"event_date = '{yesterday}'"
    )

Anti-Pattern 2: VACUUM с минимальным retention

-- НЕПРАВИЛЬНО: удаляем файлы сразу после OPTIMIZE
-- Это ломает time travel, concurrent транзакции и streaming checkpoints
SET spark.databricks.delta.retentionDurationCheck.enabled = false;
VACUUM events RETAIN 0 HOURS;  -- НИКОГДА не делайте это в production!

-- ПРАВИЛЬНО: минимум 7 дней retention
VACUUM events RETAIN 168 HOURS;

Anti-Pattern 3: coalesce(1) вместо компакции

# НЕПРАВИЛЬНО: один файл на партицию
# При большом объёме данных один Executor пишет весь раздел
# Время записи возрастает, параллелизм потерян
df.write.partitionBy("event_date").coalesce(1).parquet(path)

# ПРАВИЛЬНО: оптимальное число файлов
target = calculate_target_file_count(total_bytes, target_file_size_mb=256)
df.write.partitionBy("event_date").repartition(target).parquet(path)

Anti-Pattern 4: компакция во время пиковых запросов

# НЕПРАВИЛЬНО: компактировать в рабочее время
# OPTIMIZE требует много I/O, конкурирует с аналитическими запросами
schedule_interval = "0 9 * * 1-5"  # 09:00 в рабочие дни - плохо!

# ПРАВИЛЬНО: компактировать ночью или в выходные
schedule_interval = "0 2 * * *"   # 02:00 UTC ежедневно
# Или по выходным для исторических данных:
schedule_interval = "0 3 * * 0"   # воскресенье 03:00

Anti-Pattern 5: игнорировать файлы-сироты (Iceberg)

# НЕПРАВИЛЬНО: только expire_snapshots без remove_orphan_files
# После expire_snapshots файлы из старых snapshot'ов становятся "сиротами"
# Они физически существуют на S3, но Iceberg их не видит
# Это создаёт скрытые расходы на хранение

# Только это:
spark.sql("CALL iceberg.system.expire_snapshots(...)")
# Забыли:
# spark.sql("CALL iceberg.system.remove_orphan_files(...)")

# ПРАВИЛЬНО: всегда запускать оба шага
iceberg_maintenance(spark, ...)  # expire_snapshots + remove_orphan_files

Мониторинг и метрики компакции

from pyspark.sql import SparkSession


def audit_table_health(
    spark: SparkSession,
    table_path: str,
    format: str = "delta"
) -> dict:
    """
    Собирает метрики здоровья таблицы для принятия решения о компакции.
    """
    if format == "delta":
        dt = DeltaTable.forPath(spark, table_path)
        detail = dt.detail().first()

        file_count = detail["numFiles"]
        total_bytes = detail["sizeInBytes"]
        avg_size_mb = (total_bytes / file_count / 1024**2) if file_count > 0 else 0

        # История операций (последние 10)
        history = dt.history(10).select(
            "version", "timestamp", "operation", "operationMetrics"
        )

    elif format == "parquet":
        # Для обычного Parquet - считаем файлы через Hadoop API
        sc = spark.sparkContext
        hadoop_conf = sc._jvm.org.apache.hadoop.conf.Configuration()
        fs = sc._jvm.org.apache.hadoop.fs.FileSystem.get(
            sc._jvm.java.net.URI.create(table_path), hadoop_conf
        )
        status_list = fs.listStatus(
            sc._jvm.org.apache.hadoop.fs.Path(table_path)
        )
        file_sizes = [s.getLen() for s in status_list if not s.isDirectory()]
        file_count = len(file_sizes)
        total_bytes = sum(file_sizes)
        avg_size_mb = (total_bytes / file_count / 1024**2) if file_count > 0 else 0
        history = None
    else:
        raise ValueError(f"Неизвестный формат: {format}")

    # Классификация здоровья
    if avg_size_mb < 10:
        health = "CRITICAL"
        action = "Немедленная компакция!"
    elif avg_size_mb < 64:
        health = "WARNING"
        action = "Компакция рекомендована"
    elif avg_size_mb < 512:
        health = "OK"
        action = "Компакция не нужна"
    else:
        health = "LARGE_FILES"
        action = "Файлы слишком большие для оптимального параллелизма"

    metrics = {
        "path": table_path,
        "format": format,
        "file_count": file_count,
        "total_size_gb": total_bytes / 1024**3,
        "avg_size_mb": avg_size_mb,
        "health": health,
        "action": action
    }

    # Распределение размеров файлов
    if format == "delta" and file_count > 0:
        # Через Delta log можно получить distribution
        files_df = spark.read.format("delta").load(table_path) \
            ._jdf.queryExecution().analyzed()
        # В реальном коде: dt.detail() имеет sizeDistribution в новых версиях Delta

    return metrics


def print_audit_report(metrics: dict) -> None:
    """Форматированный вывод отчёта."""
    status_icon = {
        "CRITICAL": "🔴",
        "WARNING":  "🟡",
        "OK":       "🟢",
        "LARGE_FILES": "🔵"
    }.get(metrics["health"], "⚪")

    print(f"\n{'='*50}")
    print(f"{status_icon} Таблица: {metrics['path']}")
    print(f"  Формат:       {metrics['format']}")
    print(f"  Файлов:       {metrics['file_count']:,}")
    print(f"  Объём:        {metrics['total_size_gb']:.2f} ГБ")
    print(f"  Средний файл: {metrics['avg_size_mb']:.1f} МБ")
    print(f"  Статус:       {metrics['health']}")
    print(f"  Рекомендация: {metrics['action']}")
    print(f"{'='*50}\n")

Полный production workflow компакции

Полноценный production pipeline компакции объединяет все рассмотренные компоненты:

from datetime import datetime, timedelta
from pyspark.sql import SparkSession


def run_production_compaction_pipeline(
    spark: SparkSession,
    tables: list[dict]
) -> None:
    """
    Production pipeline компакции для нескольких таблиц.

    tables: список словарей с описанием таблиц:
    [
        {
            "path": "s3a://bucket/events/",
            "format": "delta",
            "partition_column": "event_date",
            "zorder_columns": ["user_id", "event_type"],
            "target_file_size_mb": 256
        },
        {
            "path": "s3a://bucket/users/",
            "format": "iceberg",
            "catalog": "iceberg",
            "database": "analytics",
            "table": "users",
            "strategy": "binpack",
            "target_file_size_mb": 128
        }
    ]
    """
    yesterday = (datetime.utcnow() - timedelta(days=1)).strftime("%Y-%m-%d")
    results = []

    for table_config in tables:
        fmt = table_config["format"]
        path = table_config["path"]

        print(f"\n{'='*60}")
        print(f"Обработка: {path} ({fmt})")

        # Шаг 1: проверяем, нужна ли компакция
        metrics = audit_table_health(spark, path, fmt)
        print_audit_report(metrics)

        if metrics["health"] in ("OK",) and metrics["file_count"] < 100:
            print("Пропускаем - компакция не нужна")
            results.append({"path": path, "action": "skipped"})
            continue

        # Шаг 2: компакция по формату
        if fmt == "delta":
            partition_col = table_config.get("partition_column", "event_date")
            partition_filter = f"{partition_col} = '{yesterday}'"
            zorder_cols = table_config.get("zorder_columns")

            optimize_delta_table(
                spark, path,
                partition_filter=partition_filter,
                zorder_columns=zorder_cols
            )
            results.append({"path": path, "action": "delta_optimize"})

        elif fmt == "iceberg":
            iceberg_compact(
                spark=spark,
                catalog=table_config["catalog"],
                database=table_config["database"],
                table=table_config["table"],
                partition_filter=f"event_date = '{yesterday}'",
                strategy=table_config.get("strategy", "binpack"),
                target_file_size_mb=table_config.get("target_file_size_mb", 256)
            )
            results.append({"path": path, "action": "iceberg_compact"})

    # Шаг 3: итоговый отчёт
    print(f"\n{'='*60}")
    print("Итоговый отчёт компакции:")
    for r in results:
        print(f"  {r['path']}: {r['action']}")


# Запуск pipeline
tables = [
    {
        "path": "s3a://bucket/events/",
        "format": "delta",
        "partition_column": "event_date",
        "zorder_columns": ["user_id", "event_type"],
        "target_file_size_mb": 256
    },
    {
        "path": "s3a://bucket/sessions/",
        "format": "delta",
        "partition_column": "session_date",
        "target_file_size_mb": 128
    }
]

run_production_compaction_pipeline(spark, tables)

Production кейс: от 3 часов к 4 минутам

Ситуация. Streaming-пайплайн обрабатывал события Kafka и писал результаты в Delta Lake таблицу analytics.events. Конфигурация триггера - каждые 30 секунд. Каждый батч производил ~50 файлов (50 Spark partition'ов по умолчанию, данные небольшие). За 24 часа накапливалось:

50 файлов × 2 батча/мин × 60 мин × 24 ч = 144 000 файлов/день

Через 30 дней: 4 320 000 файлов, средний размер 18 КБ.

Симптомы:

  1. spark.read.format("delta").load(path).count() - занимал 3 часа вместо 5 секунд.
  2. Driver OOM при планировании запросов: java.lang.OutOfMemoryError: GC overhead limit exceeded - 4.3M FileStatus объектов × 200B = 860 МБ только для метаданных.
  3. Delta Log checkpoint операция (checkpoint.parquet) сама занимала 45 минут: нужно записать состояние 4.3M файлов в Parquet.

Решение, шаг за шагом:

# День 1: экстренная компакция всей таблицы через OPTIMIZE
# Запускаем на ночь (02:00), ждём 4 часа
spark.sql("""
    OPTIMIZE analytics.events
""")
# Результат: 4 320 000 файлов → 1 800 файлов (средний 90 МБ)
# count() теперь: 12 секунд (вместо 3 часов)

# День 2: настраиваем ежедневный OPTIMIZE только вчерашних партиций
spark.sql("""
    OPTIMIZE analytics.events
    WHERE event_date = '2024-01-14'
""")
# Время: 2 минуты для одной партиции

# Параллельно: фиксируем стриминг с coalesce
# Вместо 50 файлов на батч - 2 файла
def fixed_batch_handler(df: DataFrame, batch_id: int) -> None:
    row_count = df.count()
    # Предполагаем ~500 байт на строку; 64 МБ целевой файл
    target_files = max(1, row_count * 500 // (64 * 1024 * 1024))
    df.coalesce(target_files).write.format("delta").mode("append").save(path)

# День 3: VACUUM для физической очистки (файлы старше 7 дней)
spark.sql("VACUUM analytics.events RETAIN 168 HOURS")
# S3: 77 ГБ старых файлов удалено

Финальные результаты через 30 дней:

Метрика До После
Файлов в таблице 4 320 000 1 800
Средний размер файла 18 КБ 90 МБ
Время count() 3 часа 12 секунд
Driver память при планировании OOM (>4 ГБ) 80 МБ
Delta Log checkpoint 45 мин 30 сек
Хранение на S3 77 ГБ 160 МБ (сжатые данные те же)
Ежедневный OPTIMIZE 4 ч на всю таблицу 2 мин на одну партицию

Ключевой вывод: компакция - это не разовая операция «починки», а необходимая часть операционного процесса для любой системы хранения, которая получает частые записи. Без планового OPTIMIZE/rewrite_data_files даже идеально написанный пайплайн деградирует до неприемлемого состояния за несколько недель.