Compaction стратегии: Spark job, Delta OPTIMIZE, Iceberg rewrite_data_files
Глубокий разбор стратегий компакции файлов в object storage: ручной Spark job, Delta Lake OPTIMIZE/VACUUM, Iceberg rewrite_data_files/expire_snapshots, автоматическая компакция в стриминге и планирование через Airflow.
Компакция как операционная практика¶
Предыдущий урок показал, почему маленькие файлы - это катастрофа для производительности: 10 000 файлов по 13 КБ читаются в 155 раз медленнее, чем один файл 128 МБ того же объёма. Причина - Time To First Byte (TTFB) в 130 мс на каждый HTTP GET-запрос, который не зависит от размера файла.
Компакция (compaction) - это процесс объединения множества мелких файлов в меньшее количество оптимальных по размеру файлов. По сути, это операция «прочитать N маленьких файлов, записать M больших файлов» (где N >> M), выполняемая периодически или по триггеру.
Компакция - не разовая операция, а непрерывная операционная практика. Данные неизбежно фрагментируются: стриминг пишет файл каждые 30 секунд, ETL делает append без coalesce, retry-механизмы создают дублирующие файлы. Без регулярной компакции любая система хранения данных деградирует.
Цели компакции:
- Сокращение числа файлов - уменьшение числа HTTP GET-запросов при чтении
- Оптимизация размера файлов - попадание в диапазон 128–512 МБ (sweet spot для Parquet)
- Кластеризация данных - во время компакции можно отсортировать данные по часто используемым в фильтрах полям, кардинально ускорив будущие запросы
- Уменьшение расходов на 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"]
)
Ограничения ручной компакции¶
Ручная компакция - наивный подход, у которого есть серьёзные проблемы:
-
Отсутствие атомарности на S3. Hadoop rename на S3 = copy + delete. Если процесс упадёт после копирования, но до удаления - получим дубли. Если после удаления, но до завершения копирования - потеряем данные.
-
Нет защиты от concurrent читателей. Пока идёт компакция, читатели видят смесь старых и новых файлов.
-
Нет отката. Если компакция записала данные неверно, отката нет - нужно перечитывать из резервной копии.
-
Нет метаданных о размерах файлов. Нужно каждый раз сканировать весь путь через
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. Причины:
-
Concurrent readers: транзакция, начавшаяся 6 часов назад, читает файлы по состоянию на момент своего начала. Если VACUUM удалит эти файлы, транзакция упадёт с FileNotFoundException.
-
Time travel:
SELECT * FROM events TIMESTAMP AS OF '2024-01-14'требует файлов, которые существовали в тот момент. Без них - ошибка. -
Backup window: если сегодня обнаружили баг в ETL, который портит данные начиная с позавчера - 7-дневное окно позволяет откатиться к чистым данным.
-
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 КБ.
Симптомы:
spark.read.format("delta").load(path).count()- занимал 3 часа вместо 5 секунд.- Driver OOM при планировании запросов:
java.lang.OutOfMemoryError: GC overhead limit exceeded- 4.3M FileStatus объектов × 200B = 860 МБ только для метаданных. - 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 даже идеально написанный пайплайн деградирует до неприемлемого состояния за несколько недель.