Merge-on-Read: delete files, position deletes и equality deletes
Анатомия Merge-on-Read в Apache Iceberg v2: устройство delete-файлов, разница между position deletes и equality deletes, как сканирование объединяет их с данными при чтении, конфигурация write.delete.mode и инженерные критерии выбора между CoW и MoR.
Почему появился Merge-on-Read¶
Прошлый урок подробно разобрал Copy-on-Write и его главную цену - Write Amplification: даже точечное изменение одной строки заставляет движок физически переписать целый data-файл, в котором эта строка находится. Для batch-обновлений с редкой частотой это приемлемая плата за полностью «чистое» состояние таблицы. Но производственный кейс того же урока показал, что происходит, когда CoW-таблица встречается с высокочастотной streaming-нагрузкой: write amplification перестаёт быть теоретической метрикой и превращается в перегруженный кластер и каскадные сбои BI-слоя.
Merge-on-Read (MoR) - это вторая стратегия Iceberg для построчного изменения данных, и она построена на прямо противоположной философии. Вместо того чтобы немедленно переписывать data-файл целиком, MoR записывает компактный файл-инструкцию о том, какие строки нужно считать удалёнными или замещёнными, и оставляет старый data-файл физически нетронутым. Эта инструкция называется delete file, и именно её анатомия - центральная тема этого урока.
Лозунг этой стратегии можно сформулировать одной фразой: «запишем изменения рядом, а разберёмся при чтении». Если CoW решает проблему немедленно и платит за это write amplification, то MoR откладывает решение проблемы до момента, когда таблицу действительно читают, и платит за это дополнительной работой на каждом SELECT. Это не бесплатный обед и не магическое устранение проблемы предыдущего урока - это перенос стоимости с одной фазы жизненного цикла данных на другую, и инженерное решение о том, какая фаза может себе позволить эту стоимость, является сердцем этого урока.
Обратите внимание на симметрию этой диаграммы с компромиссной диаграммой прошлого урока: то, что там было «WRITE дороже / READ бесплатнее», здесь буквально перевёрнуто на «WRITE почти бесплатен / READ дороже». Ни одна из двух стратегий не «лучше» другой в абсолютном смысле - выбор между ними является чисто инженерным решением, зависящим от профиля нагрузки конкретной таблицы, и именно поэтому Iceberg позволяет настраивать режим отдельно для каждой таблицы и даже отдельно для каждого типа операции (это будет показано в разделе про конфигурацию).
Read Penalty: перенос вычислений с записи на чтение¶
Термин Read Penalty описывает именно эту перенесённую стоимость. В Copy-on-Write вся вычислительная работа по согласованию состояния строки происходит один раз, синхронно, в момент выполнения UPDATE/DELETE/MERGE - и больше никогда не повторяется. В Merge-on-Read эта работа не выполняется один раз - она выполняется заново при каждом сканировании таблицы, потому что delete-файл сам по себе не меняет содержимое data-файла, а лишь предоставляет инструкцию, которую нужно применять каждый раз.
Это различие принципиально для понимания экономики MoR: если таблицу с накопленными delete-файлами читают одним запросом раз в сутки, Read Penalty платится один раз в сутки. Если ту же таблицу читают 500 раз в сутки из дашбордов, Read Penalty умножается на 500 - и именно совокупная стоимость многократного чтения, а не стоимость одной записи, обычно определяет, оправдан ли выбор MoR для конкретной таблицы. Это станет ключевым количественным аргументом в разделе про деградацию чтения практического блока.
Merge-on-Read как LSM-дерево: аналогия с журналом изменений¶
Инженерам, знакомым с LSM-деревьями (Log-Structured Merge-tree) - структурой данных, лежащей в основе таких систем, как Cassandra, RocksDB или HBase - архитектура Merge-on-Read будет интуитивно понятна, потому что Iceberg реализует очень похожий принцип на уровне табличного формата, а не движка хранения отдельной БД.
В классическом LSM-дереве запись никогда не модифицирует существующие отсортированные файлы (SSTable) на диске. Вместо этого новая запись (включая «логическое удаление», помечаемое специальным значением tombstone) добавляется в виде нового, независимого сегмента. Чтение объединяет («merge») все релевантные сегменты на лету, начиная с самых новых, чтобы определить актуальное значение каждого ключа. Периодическая операция compaction в фоне сливает несколько старых сегментов в один, физически применяя tombstone'ы и сокращая число сегментов, которые нужно объединять при следующем чтении.
Аналогия не случайна и не поверхностна: delete file Iceberg - это, по существу, табличный аналог tombstone-записи LSM-дерева, а будущая операция rewrite_data_files (детально разбираемая в десятом уроке модуля «Table Maintenance») - прямой аналог LSM-compaction. Разница лишь в масштабе и единице работы: LSM-дерево оперирует отдельными ключами внутри одного локального хранилища, а Iceberg оперирует целыми файлами в распределённом object storage, и поэтому единицей tombstone'а становится не отдельная запись в памяти, а целый файл метаданных - delete file, - который будет разобран в следующем разделе.
Эта аналогия одновременно объясняет, почему Merge-on-Read не является чем-то экзотическим или специфичным именно для Iceberg: это хорошо изученный, десятилетиями проверенный паттерн инженерии данных, адаптированный к специфике колоночного хранения в object storage, а не изобретённый с нуля.
Анатомия Delete File: новый тип файла в Iceberg v2¶
Первый урок модуля анонсировал, что Table Format Spec v2 вводит delete files как механизм Merge-on-Read, а второй урок детально разобрал sequence number - механизм, который делает применение delete-файлов корректным без какой-либо блокировки между writer'ами. Настало время собрать эти две темы в полную картину.
Delete file - это физический файл (обычно в формате Parquet, хотя спецификация допускает Avro и ORC), который не содержит «полезных» данных пользователя в привычном смысле, а содержит исключительно информацию о том, какие строки в существующих data-файлах должны считаться удалёнными или устаревшими на момент чтения. Delete file существует ровно по тем же физическим законам, что и любой Parquet-файл, разобранным в прошлом уроке (магическое число, row groups, footer) - но его логическая роль в табличной модели принципиально иная: он не описывает состояние мира, а описывает изменение этого состояния.
Спецификация Iceberg v2 выделяет ровно два типа delete-файлов, различающихся по содержимому поля content в записи манифеста:
POSITION_DELETES(значение1) - position deletes, удаление по точным координатам строки;EQUALITY_DELETES(значение2) - equality deletes, удаление по значению ключевых колонок.
Для сравнения, обычный data-файл имеет content = 0 (DATA). Эта классификация - не деталь реализации, а фундаментальное архитектурное решение, которое отражается на каждом уровне дерева метаданных, разобранного во втором уроке модуля.
Манифесты v2: разделение Data и Delete файлов¶
Второй урок модуля показал, что каждая запись манифеста (manifest entry) описывает один физический файл и имеет поле status (ADDED/EXISTING/DELETED). В v2 к этому добавляется ограничение на уровне самого манифест-файла: один манифест-файл никогда не смешивает записи о data-файлах и записи о delete-файлах - у каждого манифеста есть собственное поле content, равное либо DATA, либо DELETES, и это поле фиксировано для всего файла целиком.
Такое разделение - инженерно осмысленное решение, а не произвольное усложнение спецификации: третий урок модуля показал, что движок выполняет Scan Planning через последовательное чтение манифест-листа и манифестов, и если для конкретного запроса нужны только границы партиций и статистика data-файлов (например, для COUNT(*) без delete-файлов в этой партиции), разделение позволяет вообще не открывать delete-манифесты. На практике почти каждый реальный scan v2-таблицы открывает оба типа манифестов, но архитектурно чистое разделение упрощает как реализацию читателя, так и инкрементальную эволюцию формата - например, добавление третьего типа содержимого в будущих версиях спецификации не потребовало бы переписывания существующих data-манифестов.
Manifest list при этом, как и в случае с обычными data-файлами (третий урок модуля, поля added_data_files_count/deleted_data_files_count), хранит для каждого манифеста сводные счётчики, специфичные для delete-файлов - например, число добавленных и удалённых delete-файлов в этом манифесте. Это позволяет движку на самом верхнем уровне дерева метаданных понять, есть ли в принципе delete-файлы, релевантные для текущего снапшота, прежде чем спускаться глубже.
Связь через Sequence Number: напоминание и уточнение¶
Второй урок модуля сформулировал правило применимости delete-файла к data-файлу: delete-файл применяется к data-файлу, если sequence_number delete-файла больше или равен sequence_number данных в этом data-файле. Это правило - не побочная деталь, а единственный механизм, который делает Merge-on-Read корректным без какой-либо координации между параллельными writer'ами.
Представим ситуацию без sequence number: два процесса параллельно коммитят изменения - один пишет новый data-файл с заказами, другой удаляет старые заказы через delete-файл. Если бы порядок применения определялся, скажем, временем создания файла (а не монотонным счётчиком), рассинхронизация часов между узлами кластера могла бы привести к тому, что delete-файл «удалит» строки из data-файла, записанного после него - логическая ошибка, которая исказила бы историю изменений. Sequence number, выдаваемый централизованно catalog'ом в момент атомарного commit'а (второй урок модуля), гарантированно монотонен и устраняет этот класс ошибок полностью.
# Псевдо-модель применимости delete-файла к data-файлу,
# напоминание правила из второго урока модуля
def delete_file_applies_to(delete_file: dict, data_file: dict) -> bool:
return delete_file["sequence_number"] >= data_file["sequence_number"]
data_file = {"path": "orders_001.parquet", "sequence_number": 5}
delete_a = {"path": "delete_pos_001.parquet", "sequence_number": 6} # применяется
delete_b = {"path": "delete_pos_002.parquet", "sequence_number": 3} # НЕ применяется
print(delete_file_applies_to(delete_a, data_file)) # True
print(delete_file_applies_to(delete_b, data_file)) # False
Важное уточнение, не разбиравшееся во втором уроке: для delete-файлов помимо sequence_number (логический номер снапшота, в котором запись «появилась») существует также file_sequence_number (номер снапшота, в котором физический файл фактически был записан). Для подавляющего большинства delete-файлов эти два числа совпадают - расхождение возникает только в редких операциях типа cherry-pick снапшота из другой ветки истории, когда логическое появление записи и физическая запись файла происходят в разные моменты. Для целей этого урока достаточно считать оба числа равными - но важно знать, что спецификация различает их явно, и любой code, работающий с манифестами на низком уровне, должен использовать sequence_number для определения применимости, а не file_sequence_number.
Position Deletes: удаление по точным координатам строки¶
Position delete - наиболее распространённый и, как будет показано в разделе про конфигурацию, фактически единственный тип delete-файла, который может сгенерировать чистый Spark SQL DELETE/UPDATE/MERGE без подключения внешних потоковых систем. Принцип его работы предельно конкретен: вместо описания условия, которому соответствуют удаляемые строки, position delete описывает их точное физическое местоположение.
Схема Position Delete File¶
Файл position delete - это Parquet-файл с минимальной, строго специфицированной схемой из двух обязательных колонок:
| Колонка | Тип | Назначение |
|---|---|---|
file_path |
string |
Полный путь к data-файлу, содержащему удаляемую строку |
pos |
long |
Порядковый номер строки внутри file_path (отсчёт с нуля, в порядке физической записи) |
row (опционально) |
struct |
Полное содержимое удаляемой строки - опциональное поле для ускорения некоторых сценариев чтения |
Опциональное поле row заслуживает отдельного объяснения: спецификация v2 позволяет writer'у дополнительно записать в delete-файл не только координаты, но и полное содержимое удаляемой строки на момент удаления. Это избыточно с точки зрения чистой логики (для определения «удалена или нет» достаточно file_path и pos), но может ускорить специфические сценарии, такие как построение change data feed («что именно было удалено», а не только «что-то было удалено») - подобная возможность станет более полезной в контексте MERGE INTO (восьмой урок модуля), но базовый принцип удаления work не требует этого поля.
Обратите внимание, что один position delete file может ссылаться на строки из нескольких разных data-файлов одновременно (на диаграмме - orders_001.parquet и orders_002.parquet) - это типично для операций, затрагивающих строки, физически разбросанные по нескольким файлам, например для DELETE по непартиционирующему условию.
Жизненный цикл записи: почему нужен Scan, но дешевле, чем в CoW¶
Чтобы записать position delete, движок обязан точно знать координаты (file_path, pos) каждой удаляемой строки - а значит, обязан сначала прочитать candidate-файлы и вычислить позиции строк, удовлетворяющих условию WHERE/ON. На первый взгляд это похоже на Read-Modify-Write из прошлого урока - но решающее отличие в том, что записывается на этом шаге.
Сравнение с прошлым уроком здесь принципиально: и CoW, и position delete обязаны прочитать candidate-файлы целиком, чтобы определить, какие строки подходят под условие - стадия Scan Planning и стадия чтения candidate-файлов идентичны для обеих стратегий. Разница начинается только на шаге записи: CoW записывает заново весь файл (минус удалённые строки), а position delete записывает только список координат - объём этой записи пропорционален числу удаляемых строк, а не размеру файла, в котором они находятся.
Если в файле размером 128 МБ удаляется 50 строк, position delete file будет весить считанные килобайты (две колонки фиксированного размера на 50 записей плюс служебные накладные расходы Parquet), тогда как CoW переписал бы все 128 МБ. Это - прямое объяснение того, почему position delete устраняет именно ту часть Write Amplification, которая связана с записью неизменённых строк, но не устраняет стадию чтения candidate-файлов, которая в обоих случаях остаётся одинаково дорогой.
Поведение при чтении: Row-level Join в памяти executor'а¶
При сканировании data-файла, к которому применим хотя бы один position delete file (по правилу sequence number из предыдущего раздела), executor выполняет дополнительный шаг перед тем, как вернуть строки выше по конвейеру выполнения: он строит индекс позиций, подлежащих исключению (PositionDeleteIndex, концептуально - отсортированный набор или битовую карту номеров строк), читая все применимые position delete files для этого конкретного data-файла, и затем фильтрует поток строк data-файла по этому индексу.
# Концептуальная модель того, что делает executor при чтении
# data-файла с применимыми position delete files
def apply_position_deletes(data_file_rows, applicable_delete_files):
deleted_positions = set()
for delete_file in applicable_delete_files:
for entry in read_parquet(delete_file):
if entry["file_path"] == data_file_rows.path:
deleted_positions.add(entry["pos"])
for pos, row in enumerate(data_file_rows):
if pos not in deleted_positions:
yield row # строка не помечена как удалённая - возвращаем её
Поскольку pos хранится в порядке физической записи строк, а Parquet читается row-group за row-group в детерминированном порядке, executor может применять этот фильтр потоково, без необходимости держать в памяти весь data-файл сразу - что делает накладные расходы на чтение одного применимого position delete file относительно небольшими по сравнению с полным сканированием data-файла. Накладные расходы становятся заметными именно при накоплении множества delete-файлов на один data-файл - этот эффект будет измерен в Кейсе 3 практического блока.
Equality Deletes: удаление по значению ключа, а не по позиции¶
Position deletes требуют точного знания позиции строки в конкретном файле - а значит, требуют как минимум один scan существующих данных. Но что делать, если позиция неизвестна и сканировать существующие данные нежелательно или невозможно - например, когда изменения приходят потоком из внешней CDC-системы, которая знает только бизнес-ключ изменившейся записи (customer_id = 500), но ничего не знает о внутренней физической организации Iceberg-таблицы?
Для этого случая существует equality delete - delete-файл, который вместо координат хранит условие, описывающее удаляемые строки через значения одной или нескольких колонок.
Схема Equality Delete File¶
Файл equality delete - это Parquet-файл, схема которого состоит из произвольного подмножества колонок исходной таблицы, объявленных как ключ удаления (equality field ids - набор field id, заимствованных из той же системы стабильных идентификаторов колонок, что разбиралась в уроке про partition evolution). Содержимое файла - это конкретные значения этого ключа, подлежащие удалению.
| Колонка delete-файла | Значение в примере | Смысл |
|---|---|---|
customer_id (equality field) |
500 |
Удалить все строки с customer_id = 500 |
customer_id (equality field) |
742 |
Удалить все строки с customer_id = 742 |
Принципиальное архитектурное отличие от position delete видно прямо на диаграмме: equality delete file не знает и не обязан знать, в каких конкретно data-файлах физически находятся строки с customer_id = 500 - он формулирует условие «удалить все совпадения», а обязанность найти эти совпадения целиком перекладывается на момент чтения.
Blind Write: запись со скоростью пули¶
Эта особенность даёт equality delete его главное преимущество - возможность blind write (запись «вслепую»): writer может создать равноценный delete-файл, вообще не читая ни байта существующих данных таблицы. Достаточно знать значение ключа, которое нужно удалить - и больше ничего.
Это именно то свойство, которое делает equality deletes фундаментом потоковой интеграции Iceberg: источник CDC (например, Debezium, читающий write-ahead log PostgreSQL) сообщает только что «строка с таким первичным ключом удалена», и writer может немедленно зафиксировать это как новый снапшот без какой-либо задержки на чтение текущего состояния многотерабайтной таблицы. Задержка commit'а equality delete не зависит от размера таблицы вообще - в отличие от position delete, где чтение candidate-файлов (хотя и дешевле полного CoW-rewrite) всё же растёт с объёмом данных, попадающих под условие.
Скрытая угроза: fan-out join при чтении¶
Цена этой скорости записи проявляется ровно в противоположной точке жизненного цикла - при чтении. Поскольку equality delete file не содержит координат, единственный способ применить его - проверить каждую потенциально релевантную строку каждого применимого data-файла на совпадение с условием удаления. Это и есть полноценный join: для каждой строки каждого применимого data-файла движок должен сравнить значение колонки customer_id со списком значений в equality delete file.
Партиционирование таблицы (transform-функции, разобранные в четвёртом уроке) частично смягчает эту угрозу: equality delete file, как и любой другой файл данных, привязан к конкретной партиции, и поэтому применяется только к data-файлам той же партиции, а не ко всей таблице. Но внутри одной партиции, если она содержит десятки или сотни data-файлов, единственный equality delete file заставляет движок выполнить join этого delete-файла с каждым из них - и именно это накопление становится опасным, если компакция (десятый урок) не выполняется регулярно. Этот эффект будет измерен в Кейсе 3 практического блока этого урока.
Position Deletes vs Equality Deletes: сравнительная таблица¶
| Критерий | Position Deletes | Equality Deletes |
|---|---|---|
| Что хранится | (file_path, pos) - точные координаты |
Значения ключевых колонок (equality_ids) |
| Нужен ли Scan перед записью | Да - нужно найти позиции совпадающих строк | Нет - blind write, можно писать без чтения таблицы |
| Стоимость записи | Пропорциональна числу удаляемых строк + стоимость scan candidate-файлов | Минимальна и не зависит от объёма таблицы |
| Стоимость применения при чтении | Дешёвая - прямое сопоставление по точным координатам | Дорогая - требует join со всеми применимыми data-файлами |
| Кто обычно генерирует | Spark SQL DELETE/UPDATE/MERGE (Iceberg Spark extensions) |
Потоковые writer'ы (Flink CDC sink, Kafka Connect) без доступа к scan |
| Идеальный сценарий | Batch-операции внутри одного движка, знающего позиции строк | Высокочастотный CDC, blind upserts по первичному ключу |
| Главный риск | Накопление множества delete-файлов на один data file | Fan-out join со всеми data-файлами партиции при отсутствии компакции |
Это сравнение - не абстрактная теория, а прямое следствие архитектурных решений, разобранных в предыдущих двух разделах: position delete платит небольшую цену при записи (scan, как у CoW, но без полной перезаписи) за дешёвое чтение; equality delete платит почти нулевую цену при записи (blind write) за потенциально дорогое чтение, которое дополнительно деградирует со временем без обслуживания.
Полный путь SELECT под Merge-on-Read: расширение Этапа 3¶
Третий урок модуля детально разобрал четырёхэтапный алгоритм Scan Planning и отдельно выделил «Этап 3.1» - дополнительный подшаг для v2-таблиц, на котором для каждого data-файла, прошедшего pruning по статистике, движок ищет все применимые delete-файлы по правилу sequence number. Теперь, когда устройство position и equality deletes разобрано детально, можно увидеть полную картину того, что происходит дальше - на самом исполнении, а не только на этапе планирования.
Главное наблюдение этой диаграммы: один и тот же data-файл может иметь несколько применимых delete-файлов одновременно, в том числе разных типов - один position и один equality, как на диаграмме выше. Executor обязан применить их все, причём именно в порядке sequence_number, потому что более поздний delete-файл логически отражает более актуальное состояние удалений. Это - прямое объяснение того, почему стоимость Этапа 3.1, отмеченная в третьем уроке как «не катастрофическая, но заметная», растёт не линейно с числом delete-файлов, а скорее с произведением «число data-файлов × среднее число применимых delete-файлов на data-файл» - и именно это произведение является количественной мерой Read Penalty, введённого в первом разделе этого урока.
Read Penalty как растущая, а не постоянная величина¶
Важное отличие Read Penalty от Write Amplification прошлого урока: Write Amplification CoW-операции - это разовая величина, фиксируемая в момент commit'а и не меняющаяся постфактум. Read Penalty Merge-on-Read-таблицы, напротив, накапливается со временем: каждая новая MoR-операция добавляет ещё один delete-файл, потенциально применимый к уже существующим data-файлам, и стоимость каждого следующего SELECT от этой же таблицы постепенно растёт - до тех пор, пока не будет выполнена компакция (десятый урок модуля).
# Концептуальная (учебная) модель роста стоимости чтения партиции
# с накапливающимися delete-файлами без компакции
def estimate_read_penalty(data_files: int, delete_files_per_data_file: list[int]) -> dict:
total_merge_operations = sum(delete_files_per_data_file)
avg_deletes_per_file = total_merge_operations / data_files if data_files else 0
return {
"data_files": data_files,
"total_merge_operations": total_merge_operations,
"avg_deletes_per_data_file": round(avg_deletes_per_file, 2),
}
# День 1: партиция только что создана, delete-файлов нет
day_1 = estimate_read_penalty(data_files=40, delete_files_per_data_file=[0] * 40)
# День 30: 30 дней point-DELETE по той же партиции без компакции,
# каждый день затрагивает в среднем 3 разных data-файла
day_30 = estimate_read_penalty(data_files=40, delete_files_per_data_file=[3] * 40)
print(day_1) # {'data_files': 40, 'total_merge_operations': 0, 'avg_deletes_per_data_file': 0.0}
print(day_30) # {'data_files': 40, 'total_merge_operations': 120, 'avg_deletes_per_data_file': 3.0}
Эта модель учебная и сильно упрощённая (в реальности распределение delete-файлов по data-файлам неравномерно), но её вывод верен качественно: без компакции avg_deletes_per_data_file монотонно растёт, и каждый последующий SELECT от этой партиции выполняет всё больше операций merge на executor'ах - притом что объём данных в партиции мог не измениться вообще. Это - количественное обоснование того, почему «MoR без compaction-расписания» является эксплуатационным антипаттерном, а не просто теоретическим риском, и почему десятый урок модуля - не опциональное продолжение, а необходимая часть production-эксплуатации любой MoR-таблицы.
Формальная оценка стоимости merge: position vs equality¶
Качественная модель выше показывает направление роста стоимости, но для position и equality deletes форма этой стоимости различается, и это различие можно зафиксировать формально, опираясь на устройство каждого типа, разобранное в соответствующих разделах.
Для position deletes стоимость применения к одному data-файлу пропорциональна числу позиций, которые нужно проверить, - то есть числу строк, помеченных как удалённые во всех применимых delete-файлах для этого data-файла. Поскольку проверка одной позиции - это операция за O(1) относительно отсортированного индекса позиций (раздел про Row-level Join в памяти), общая стоимость растёт линейно с суммарным числом удалённых строк:
def position_delete_merge_cost(applicable_delete_files: list[dict]) -> int:
"""Условная единица стоимости: одна проверка позиции = 1 единица."""
return sum(f["record_count"] for f in applicable_delete_files)
# Кейс 3: 25 delete-файлов, каждый удаляет ~40 строк
files = [{"record_count": 40} for _ in range(25)]
print(position_delete_merge_cost(files)) # 1000 условных единиц на один data-файл
Для equality deletes ситуация принципиально иная: стоимость применения одного equality delete file к data-файлу не зависит от числа значений ключа внутри него (это - join, и его стоимость определяется числом строк data-файла, которые нужно проверить, а не числом значений в delete-файле), а главный множитель - это число применимых equality delete files, каждый из которых требует отдельного прохода join по data-файлу:
def equality_delete_merge_cost(data_file_rows: int, applicable_eq_delete_files: int) -> int:
"""Условная единица стоимости: один join data-файла с одним eq-delete = data_file_rows проверок."""
return data_file_rows * applicable_eq_delete_files
# Производственный кейс: data-файл на 50 000 строк, 180 применимых equality delete files
print(equality_delete_merge_cost(50_000, 180)) # 9 000 000 условных единиц на один data-файл
Разница на три порядка между двумя вычисленными значениями (1 000 против 9 000 000 условных единиц) - не случайность, а прямое следствие архитектурного различия: position delete «знает», какую конкретно строку проверять, и платит только за неё; equality delete не знает позиции и обязан проверить каждую строку data-файла на совпадение с условием, причём отдельно для каждого применимого delete-файла. Именно это формальное соотношение объясняет, почему производственный кейс этого урока (equality deletes без компакции) деградировал настолько сильнее, чем гипотетическое накопление такого же числа position delete files привело бы в Кейсе 3.
Конфигурирование Merge-on-Read¶
Прошлый урок представил три свойства таблицы - write.update.mode, write.delete.mode, write.merge.mode - каждое из которых принимает значение copy-on-write или merge-on-read и настраивается независимо для каждого типа операции. Этот урок возвращается к тому же механизму конфигурации, но уже фокусируясь на значении merge-on-read.
CREATE TABLE prod.db.orders (
order_id BIGINT,
customer_id BIGINT,
status STRING,
amount DECIMAL(10, 2),
event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(event_ts))
TBLPROPERTIES (
'format-version' = '2',
'write.update.mode' = 'merge-on-read',
'write.delete.mode' = 'merge-on-read',
'write.merge.mode' = 'merge-on-read'
);
Дополнительно к этим трём свойствам, специфичным для MoR, существует свойство, определяющее формат файла, в котором записываются delete-файлы - по аналогии с write.format.default для data-файлов:
ALTER TABLE prod.db.orders SET TBLPROPERTIES (
'write.delete.format.default' = 'parquet'
);
Значение по умолчанию для большинства актуальных версий Apache Iceberg при записи через Spark - parquet, что обеспечивает единообразие со всеми остальными файлами таблицы (включая поддержку статистики min/max и предикатного pruning средствами самого Parquet). Некоторые потоковые writer'ы (в частности, более старые версии Flink-коннектора) исторически писали delete-файлы в формате Avro - построчном, без колоночного сжатия, что было разумным компромиссом для очень маленьких, часто создаваемых файлов, где накладные расходы колоночного формата не оправдывали себя. На практике для self-hosted стенда, описанного в этом модуле (Spark + JDBC Catalog + MinIO), рекомендуется явно зафиксировать parquet для предсказуемости и единообразия инструментов анализа метаданных.
Дерево решений: выбор режима по профилю нагрузки¶
Прежде чем переходить к табличной матрице, полезно представить тот же выбор в виде последовательности вопросов, которые инженер задаёт себе при проектировании конкретной таблицы - именно в таком порядке решение принимается на практике:
Обратите внимание, что узел Q3 фактически не является свободным выбором инженера - как было показано в Кейсе 1 и Кейсе 2 практического блока, Spark всегда располагает позициями к моменту записи delete-файла, и поэтому ветка «нет» из Q3 на практике недостижима для Spark SQL - выбор между position и equality deletes определяется не настройкой, а тем, какой именно writer выполняет операцию.
Сравнительная матрица: CoW vs MoR (Position) vs MoR (Equality)¶
| Паттерн нагрузки | Рекомендованный режим | Обоснование |
|---|---|---|
| Редкие batch-корректировки (раз в день/неделю) | copy-on-write |
Read Penalty был бы постоянным грузом на read-heavy BI-слой; CoW платит цену один раз при записи (прошлый урок) |
Частые точечные UPDATE/DELETE из одного batch-движка (Spark) |
merge-on-read (position deletes) |
Spark знает позиции строк бесплатно в рамках своего собственного scan; запись дешевле CoW, чтение остаётся относительно предсказуемым |
| High-frequency CDC от внешнего источника (Debezium/Kafka) | merge-on-read (equality deletes, через потоковый writer) |
Источник знает только бизнес-ключ, не положение строки; blind write устраняет задержку, пропорциональную размеру таблицы |
| Near real-time streaming с десятками commits в минуту | merge-on-read + обязательное регулярное rewrite_data_files (десятый урок) |
Без компакции Read Penalty растёт без ограничений (см. модель выше); компакция - не опция, а обязательная часть архитектуры |
| Batch-обновление, выполняемое самим Spark, но читаемое сразу после записи в том же pipeline | merge-on-read (position deletes) с компакцией в конце pipeline |
Промежуточные чтения внутри одного pipeline допускают небольшой Read Penalty, если финальный шаг материализует delete-файлы перед публикацией |
GDPR/право на забвение, единичные point-DELETE по customer_id |
copy-on-write |
Та же рекомендация, что и в прошлом уроке: редкая, юридически значимая операция не должна оставлять следов в виде delete-файлов дольше, чем необходимо |
Важно зафиксировать: эта матрица - не догма, а отправная точка для инженерного решения, которое всегда должно проверяться измерением на реальной нагрузке, как было показано на примере прошлого урока (производственный кейс с CoW-таблицей, столкнувшейся со streaming-нагрузкой). Симметричная ошибка - использовать merge-on-read для read-heavy таблицы с редкими изменениями - тоже возможна и будет разобрана в разделе про типичные заблуждения.
Use cases и антипаттерны Merge-on-Read¶
Идеальные случаи применения:
- CDC-репликация из OLTP-систем - PostgreSQL/MySQL через Debezium, где источник поставляет события по первичному ключу без знания о физической организации Iceberg-таблицы; equality deletes устраняют необходимость scan при каждом изменении.
- Bronze/Silver-слои с высокочастотным upsert (медальонная архитектура) - таблицы, в которые данные льются непрерывным потоком и которые читаются относительно редко (для построения Silver/Gold из Bronze), что делает Read Penalty приемлемым при условии регулярной компакции.
- Таблицы с частыми point-DELETE по неиндексируемому условию в batch-режиме Spark - position deletes избегают дорогого CoW rewrite, оставаясь при этом относительно дешёвыми для чтения.
Антипаттерны:
- Read-heavy BI-витрины и Gold-слой без расписания компакции - именно тот случай, где прошлый урок рекомендовал CoW; если же выбран MoR, отсутствие компакции напрямую транслируется в деградацию дашбордов.
- Equality deletes без plan на компакцию для широких партиций - fan-out join, разобранный в соответствующем разделе, делает эту комбинацию особенно опасной при сотнях data-файлов на партицию.
- MoR «на всякий случай», без измеренного профиля нагрузки - симметрично антипаттерну прошлого урока («CoW по умолчанию»); выбор режима должен опираться на измерение, а не на интуицию или настройки по умолчанию.
Практический демо-блок: инспекция delete-файлов в PySpark¶
Стенд использует ту же self-hosted инфраструктуру, что и весь модуль: JDBC Catalog поверх PostgreSQL для метаданных и S3FileIO поверх MinIO для физических файлов (полная конфигурация SparkSession приведена в первом уроке модуля и здесь не повторяется). Все примеры ниже выполняются в spark-shell/pyspark с подключёнными Iceberg Spark extensions.
Кейс 0: создание MoR-таблицы и базовое наполнение¶
spark.sql("""
CREATE TABLE lakehouse.analytics.orders (
order_id BIGINT,
customer_id BIGINT,
status STRING,
amount DECIMAL(10, 2),
event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(event_ts))
TBLPROPERTIES (
'format-version' = '2',
'write.delete.mode' = 'merge-on-read',
'write.update.mode' = 'merge-on-read',
'write.merge.mode' = 'merge-on-read'
)
""")
import random
from datetime import datetime, timedelta
base_ts = datetime(2024, 1, 1)
statuses = ["created", "paid", "shipped", "cancelled", "refunded"]
rows = [
(
order_id,
random.randint(1, 50_000),
random.choice(statuses),
round(random.uniform(5, 500), 2),
base_ts + timedelta(days=order_id % 30, seconds=order_id),
)
for order_id in range(1, 1_000_001)
]
df = spark.createDataFrame(
rows,
schema="order_id BIGINT, customer_id BIGINT, status STRING, amount DECIMAL(10,2), event_ts TIMESTAMP",
)
df.writeTo("lakehouse.analytics.orders").append()
spark.sql("""
SELECT count(*) AS total_files, sum(file_size_in_bytes) AS total_bytes
FROM lakehouse.analytics.orders.files
""").show()
# +-----------+------------+
# |total_files| total_bytes|
# +-----------+------------+
# | 30| 38734218 |
# +-----------+------------+
Тридцать файлов - ровно по одному на каждый день из тридцатидневного диапазона event_ts, использованного для генерации данных (партиционирование days(event_ts), четвёртый урок модуля). На этом шаге .delete_files (системная таблица, специфичная для v2, которая появится в следующем кейсе) пуста - в таблице существуют только обычные data-файлы.
Кейс 1: генерация и инспекция Position Deletes¶
Выполним DELETE по условию, которое не совпадает с ключом партиционирования (status, а не event_ts) - то есть условие, удовлетворяющее строки, физически разбросанные внутри каждой партиции:
spark.sql("""
DELETE FROM lakehouse.analytics.orders
WHERE status = 'cancelled' AND customer_id = 12345
""")
Поскольку таблица сконфигурирована с write.delete.mode = 'merge-on-read', и Spark выполняет эту операцию через собственный scan (а значит, точно знает позиции совпадающих строк), результатом будет position delete file, а не переписанный data-файл. Проверим это через системную таблицу .delete_files, появляющуюся в Iceberg Spark extensions именно для v2-таблиц:
spark.sql("""
SELECT content, file_path, file_format, record_count, file_size_in_bytes
FROM lakehouse.analytics.orders.delete_files
""").show(truncate=False)
# +-----------------+----------------------------------------+-----------+------------+-------------------+
# |content |file_path |file_format|record_count|file_size_in_bytes |
# +-----------------+----------------------------------------+-----------+------------+-------------------+
# |POSITION_DELETES |s3://lakehouse/.../delete-0001-pos.parquet|PARQUET | 23| 4108 |
# +-----------------+----------------------------------------+-----------+------------+-------------------+
content = POSITION_DELETES подтверждает теоретическую модель: Spark SQL DELETE в MoR-режиме генерирует именно position deletes. Обратите внимание на record_count = 23 (число фактически удалённых строк, разбросанных по нескольким партициям из-за условия по customer_id) и file_size_in_bytes = 4108 - чуть больше 4 КБ для целого delete-файла, против десятков мегабайт, которые потребовал бы CoW-rewrite затронутых партиций. Это - прямое практическое подтверждение раздела про «запись со скоростью пули», но именно для position deletes: запись здесь не blind (scan всё же выполнен), но объём записанных байт пропорционален числу удалённых строк, а не размеру содержащих их файлов.
Убедимся, что исходные data-файлы остались физически нетронутыми:
spark.sql("""
SELECT snapshot_id, summary['added-data-files'] AS added_data,
summary['deleted-data-files'] AS deleted_data,
summary['added-delete-files'] AS added_delete,
summary['added-position-deletes'] AS added_pos_deletes
FROM lakehouse.analytics.orders.snapshots
ORDER BY committed_at DESC LIMIT 1
""").show()
# +-----------+----------+------------+------------+------------------+
# |snapshot_id|added_data|deleted_data|added_delete|added_pos_deletes |
# +-----------+----------+------------+------------+------------------+
# | 784213... | 0| 0| 1| 23 |
# +-----------+----------+------------+------------+------------------+
Сравните этот summary с аналогичной проверкой из прошлого урока для CoW-таблицы: там deleted-data-files и added-data-files были ненулевыми и равными числу строк целого bucket'а (тысячи записей), потому что данные физически переписывались. Здесь added-data-files = 0 и deleted-data-files = 0 - ни один data-файл не был тронут; вместо этого появилось новое поле added-delete-files = 1 и added-position-deletes = 23, отражающее точное число логически удалённых строк. Это и есть формальное, измеримое доказательство главного тезиса урока: MoR не переписывает данные, он добавляет инструкцию рядом.
Кейс 2: эмуляция equality deletes для CDC-пайплайна¶
Важная техническая оговорка, прежде чем переходить к коду: чистый Spark SQL DELETE/UPDATE/MERGE не генерирует equality deletes - именно потому, что Spark всегда выполняет собственный scan и поэтому всегда точно знает позиции совпадающих строк (это было продемонстрировано в Кейсе 1). Планировщик row-level операций Iceberg-Spark-интеграции спроектирован так, чтобы предпочитать position deletes, как более дешёвые при чтении, в любой ситуации, где позиции в принципе доступны.
Equality deletes на практике почти всегда порождает не Spark, а потоковый writer, у которого позиций нет и не может быть - например, Flink-коннектор Iceberg, работающий в upsert-режиме поверх потока CDC-событий от Debezium. Чтобы изучить структуру equality delete файла, не разворачивая полноценный Flink-кластер, обратимся к низкоуровневому Java API Iceberg напрямую через _jvm-мост PySpark - это нестандартная техника, выходящая за рамки публичного Spark SQL API, но она показывает ровно то, что внутри себя делает любой потоковый коннектор, и часто используется инженерами для тестирования или прямой интеграции:
# Концептуальная иллюстрация низкоуровневого Java API Iceberg,
# показывающая принцип работы Flink/CDC-коннектора "изнутри".
# Это не готовый к копированию production-код, а учебная демонстрация
# того, какие шаги выполняет blind write equality delete.
jvm = spark._jvm
catalog = jvm.org.apache.iceberg.spark.Spark3Util.loadIcebergTable(
spark._jsparkSession, "lakehouse.analytics.orders"
)
schema = catalog.schema()
spec = catalog.spec()
equality_field_ids = [schema.findField("customer_id").fieldId()]
# GenericAppenderFactory создаёт writer, который пишет ТОЛЬКО
# значения equality-колонок - без какого-либо обращения к data-файлам
appender_factory = jvm.org.apache.iceberg.data.GenericAppenderFactory(
schema, spec, equality_field_ids, schema.select("customer_id"), None
)
output_file = catalog.io().newOutputFile(
catalog.locationProvider().newDataLocation("delete-eq-cdc-001.parquet")
)
eq_writer = appender_factory.newEqDeleteWriter(
output_file, jvm.org.apache.iceberg.FileFormat.PARQUET, None
)
# CDC-событие "DELETE customer_id=98214" - НИКАКОГО чтения существующих данных
record = jvm.org.apache.iceberg.data.GenericRecord.create(schema.select("customer_id"))
record.setField("customer_id", 98214)
eq_writer.write(record)
eq_writer.close()
delete_file = eq_writer.toDeleteFile()
row_delta = catalog.newRowDelta()
row_delta.addDeletes(delete_file)
row_delta.commit()
Ключевой момент, который этот псевдокод иллюстрирует: между получением CDC-события и row_delta.commit() не происходит ни одного обращения к существующим data-файлам таблицы - ни Scan Planning, ни чтения candidate-файлов, ничего из того, что было обязательным шагом в Кейсе 1 для position deletes. Это и есть blind write на практике, не только в теории.
После такого commit'а (выполненного, например, реальным Flink-коннектором на production-стенде) результат полностью наблюдаем через тот же системный путь, что и в Кейсе 1:
spark.sql("""
SELECT content, file_path, record_count, file_size_in_bytes, equality_ids
FROM lakehouse.analytics.orders.delete_files
WHERE content = 'EQUALITY_DELETES'
""").show(truncate=False)
# +-----------------+-------------------------------------------+------------+-------------------+-------------+
# |content |file_path |record_count|file_size_in_bytes |equality_ids |
# +-----------------+-------------------------------------------+------------+-------------------+-------------+
# |EQUALITY_DELETES |s3://lakehouse/.../delete-eq-cdc-001.parquet| 1| 891 | [2]|
# +-----------------+-------------------------------------------+------------+-------------------+-------------+
equality_ids = [2] ссылается на стабильный field id колонки customer_id (механизм стабильных идентификаторов разобран в уроке про partition evolution) - именно по этому полю движок будет сопоставлять строки при чтении. file_size_in_bytes = 891 байт - размер, который физически не может зависеть от объёма таблицы, потому что writer не читал из неё ни строки; сравните это с 4108 байтами position delete file из Кейса 1, который пусть и не равен размеру CoW-rewrite, но всё же пропорционален числу совпавших строк (23 строки против одной).
Кейс 3: деградация чтения от накопленных delete-файлов¶
Накопим на одной и той же партиции множество мелких position delete files - типичный паттерн для таблицы, в которую систематически льются точечные DELETE/UPDATE без регулярной компакции:
import time
# Симулируем 25 дней точечных DELETE по разным customer_id
# в рамках одной и той же партиции (event_ts в январе 2024)
for day in range(25):
target_customers = random.sample(range(1, 50_000), k=40)
spark.sql(f"""
DELETE FROM lakehouse.analytics.orders
WHERE customer_id IN ({",".join(map(str, target_customers))})
AND event_ts >= '2024-01-01' AND event_ts < '2024-01-02'
""")
delete_file_count = spark.sql("""
SELECT count(*) AS n
FROM lakehouse.analytics.orders.delete_files
WHERE content = 'POSITION_DELETES'
""").collect()[0]["n"]
print(f"Накопленных position delete files: {delete_file_count}")
# Накопленных position delete files: 25
Теперь измерим время выполнения одного и того же аналитического запроса к затронутой партиции - до того, как накопились delete-файлы, измерение было бы тривиально быстрым (чистое колоночное чтение одного файла), а сейчас executor обязан применить все 25 применимых delete-файлов к этому же data-файлу:
start = time.time()
result = spark.sql("""
SELECT count(*), sum(amount)
FROM lakehouse.analytics.orders
WHERE event_ts >= '2024-01-01' AND event_ts < '2024-01-02'
""").collect()
elapsed_with_deletes = time.time() - start
print(f"SELECT с 25 применимыми delete-файлами: {elapsed_with_deletes:.2f} сек")
# SELECT с 25 применимыми delete-файлами: 4.81 сек
spark.sql("""
SELECT data_file.file_path, count(*) AS applicable_deletes
FROM lakehouse.analytics.orders.entries
WHERE data_file.content = 0
GROUP BY data_file.file_path
ORDER BY applicable_deletes DESC
LIMIT 1
""").show(truncate=False)
# +------------------------------------------+-------------------+
# |file_path |applicable_deletes |
# +------------------------------------------+-------------------+
# |s3://lakehouse/.../00001-day01.parquet | 25|
# +------------------------------------------+-------------------+
Один и тот же data-файл партиции 2024-01-01 теперь обязан проходить через 25 раундов проверки позиций при каждом сканировании - именно то поведение, которое было предсказано моделью Read Penalty в теоретической части. Важно: при сравнимом запросе сразу после Кейса 0 (до накопления delete-файлов) аналогичный SELECT по одной партиции выполнялся в районе 0.3-0.4 секунды - рост примерно в 12-15 раз для одной и той же логической партиции, без единого байта новых «полезных» данных.
Для сравнения смоделируем, во что обошлось бы то же число операций изменения, если бы они были выполнены через equality deletes (как делал бы потоковый CDC-writer из Кейса 2), а не через Spark SQL DELETE:
applicable_data_files_in_partition = spark.sql("""
SELECT count(*) AS n
FROM lakehouse.analytics.orders.files
WHERE partition.event_ts_day = '2024-01-01'
""").collect()[0]["n"]
estimated_eq_cost = equality_delete_merge_cost(
data_file_rows=31_250, # средний размер data-файла партиции
applicable_eq_delete_files=25, # то же число операций, что и в Кейсе 3
)
estimated_pos_cost = position_delete_merge_cost([{"record_count": 40}] * 25)
print(f"data-файлов в партиции: {applicable_data_files_in_partition}")
print(f"оценка стоимости merge (position): {estimated_pos_cost}")
print(f"оценка стоимости merge (equality): {estimated_eq_cost}")
# data-файлов в партиции: 1
# оценка стоимости merge (position): 1000
# оценка стоимости merge (equality): 781250
Разница почти в 800 раз для одного и того же числа операций изменения (25 точечных delete) - количественное подтверждение формальной модели из теоретической части на данных, фактически измеренных в этом же кейсе. Это не означает, что equality deletes были бы «ошибкой» в этом сценарии - источником данных был сам Spark, у которого позиции доступны бесплатно, поэтому реальный writer (Кейс 1) совершенно обоснованно выбрал position deletes. Но если бы тот же объём изменений поступал из внешнего CDC-потока без доступа к позициям, единственной практичной альтернативой были бы equality deletes - и именно тогда регулярная компакция перестаёт быть желательной практикой и становится обязательным условием эксплуатации, что и подтвердил производственный кейс этого урока.
На этом шаге курс намеренно НЕ выполняет компакцию - rewrite_data_files/rewrite_position_delete_files, которые материализовали бы все 25 delete-файлов обратно в чистый Parquet и устранили бы наблюдаемую деградацию, - это полноценная тема десятого урока модуля «Table Maintenance». Здесь важно зафиксировать сам факт деградации и её прямую связь с числом накопленных delete-файлов на партицию, измеренную через .entries.
На этом шаге курс намеренно НЕ выполняет компакцию - rewrite_data_files/rewrite_position_delete_files, которые материализовали бы все 25 delete-файлов обратно в чистый Parquet и устранили бы наблюдаемую деградацию, - это полноценная тема десятого урока модуля «Table Maintenance». Здесь важно зафиксировать сам факт деградации и её прямую связь с числом накопленных delete-файлов на партицию, измеренную через .entries.
Производственный кейс: накопление tombstone'ов на Bronze-слое¶
Команда data-платформы крупного маркетплейса построила Bronze-слой медальонной архитектуры (модуль про data modeling этого курса) на Iceberg-таблице bronze.orders_raw, принимающей CDC-поток от Debezium через Kafka Connect с частотой несколько тысяч событий в минуту - точечные INSERT/UPDATE/DELETE, отражающие изменения в исходной OLTP-базе. Таблица была сконфигурирована с write.merge.mode = 'merge-on-read' и write.delete.mode = 'merge-on-read' - архитектурно верное решение, прямо следующее из материала этого урока: equality deletes от потокового коннектора устраняют зависимость задержки commit'а от размера таблицы.
Maintenance-job, который должен был регулярно вызывать rewrite_data_files (десятый урок модуля), был запланирован, но из-за ошибки в Airflow DAG (неправильно указанная зависимость задачи) не выполнялся в течение шести недель - и это осталось незамеченным, потому что сама запись в таблицу продолжала работать без единой ошибки.
Корневая причина инцидента - не сам выбор Merge-on-Read (он был корректен для профиля нагрузки CDC-потока), а отсутствие операционного контроля за компакцией: maintenance-процедура существовала, но её фактическое выполнение никем не отслеживалось. Это прямо иллюстрирует тезис раздела про Read Penalty: стоимость MoR не статична, она накапливается, и единственным механизмом, сдерживающим этот рост, является регулярная, проверяемая компакция - а не сам факт включения режима MoR в конфигурации таблицы.
| Метрика | Неделя 1 | Неделя 5 (пик) | После rewrite_data_files |
|---|---|---|---|
| Equality delete files | 340 | 6 100 | 0 |
| Среднее число применимых deletes на data-файл | ~3 | ~180 | 0 |
Время SELECT по партиции дня |
~2 сек | ~51 сек | ~3 сек |
| Время Silver-job (полная инкрементальная загрузка) | ~4 мин | timeout (>15 мин) | ~5 мин |
Меры по итогам инцидента:
- Немедленная мера: экстренный запуск
rewrite_data_files(синтаксис и параметры - тема десятого урока) восстановил производительность чтения до уровня, близкого к исходному. - Структурное исправление: maintenance-job был переведён из Airflow DAG с неявной зависимостью в отдельный, независимо мониторящийся scheduled job с алертом по метрике
delete_file_countиз.delete_files, превышающей пороговое значение для таблицы. - Архитектурное уточнение: команда зафиксировала правило для всех будущих MoR-таблиц - maintenance-job обязателен и верифицируется при code review конфигурации таблицы, а не добавляется «потом, когда понадобится», то есть рассматривается как неотъемлемая часть архитектуры MoR-таблицы, а не как опциональная оптимизация.
Типичные заблуждения¶
«Merge-on-Read - это бесплатная замена Copy-on-Write» - неверно и прямо противоречит главному тезису урока. MoR не устраняет стоимость согласования состояния строки, а переносит её с записи на чтение - причём эта стоимость со временем растёт, в отличие от разовой стоимости CoW-операции.
«Spark SQL DELETE генерирует equality deletes, если условие написано по бизнес-ключу» - неверно. Spark всегда выполняет собственный scan перед записью delete-файла и поэтому всегда точно знает позиции совпадающих строк - Iceberg-Spark планировщик row-level операций всегда предпочитает position deletes, независимо от формы условия WHERE. Equality deletes на практике порождают потоковые writer'ы без доступа к scan, что было показано в Кейсе 2.
«Equality deletes - всегда плохой выбор, потому что они дороже при чтении» - неточно. Для своего целевого сценария (blind write от потокового источника без доступа к позициям) equality deletes - единственный практичный вариант; «дороговизна при чтении» становится проблемой только при отсутствии регулярной компакции, что и показал производственный кейс этого урока.
«Position deletes никогда не требуют scan, поэтому они всегда дешевле equality deletes при записи» - неверно: Кейс 1 и Кейс 2 практического блока напрямую сравнили размеры файлов (4108 байт против 891 байта) и показали, что position delete всё же требует чтения candidate-файлов для определения позиций, тогда как equality delete - не требует вообще. Position deletes дешевле именно при чтении, а не при записи.
«Раз delete file маленький, он не влияет на производительность» - неверно в совокупности: размер отдельного delete-файла действительно мал (килобайты против мегабайт data-файла), но Кейс 3 практического блока показал, что деградация чтения определяется не размером отдельного файла, а количеством применимых delete-файлов на data-файл, накапливающимся со временем.
«Format-version v2 автоматически означает, что таблица использует Merge-on-Read» - повтор заблуждения из прошлого урока в обратную сторону: v2 - спецификация, допускающая delete-файлы, а не обязывающая их использовать; таблица может быть v2 и работать исключительно в режиме Copy-on-Write, если так сконфигурированы write.update.mode/write.delete.mode/write.merge.mode.
«Компакция MoR-таблицы - опциональная оптимизация, можно отложить» - прямо опровергнуто производственным кейсом этого урока: для MoR-таблицы с продолжающимся потоком изменений компакция - не оптимизация, а необходимое условие сохранения приемлемой производительности чтения, и её отсутствие - вопрос времени до инцидента, а не вопрос «нужна она вообще или нет».
«Если delete-файлов мало, не важно, position они или equality» - неточно при экстраполяции на будущее: раздел про формальную оценку стоимости merge показал, что разница в стоимости между двумя типами растёт не линейно, а как минимум мультипликативно с размером data-файла, и небольшое число equality delete files на крупных партициях может уже сейчас создавать заметный Read Penalty, даже когда аналогичное число position delete files остаётся незаметным.
Мостик к следующим урокам¶
Этот урок закрыл вторую из двух стратегий построчного изменения данных в Iceberg, показав её зеркальную симметрию с Copy-on-Write: там, где CoW платит при записи и получает бесплатное чтение, MoR платит при чтении (причём по нарастающей) и получает почти бесплатную запись. Следующие уроки модуля строят на этом фундаменте напрямую.
-
Урок 8 («MERGE INTO») разбирает production-паттерн upsert для CDC-потоков на уровне SQL-синтаксиса и конфигурации
write.merge.mode, опираясь на материал position и equality deletes этого урока - именноMERGE INTOчаще всего оказывается операцией, которая на практике генерирует оба типа delete-файлов в зависимости от того, какая часть условияWHEN MATCHEDзатронута. -
Урок 9 («Time Travel») возвращается к delete-файлам в контексте истории снапшотов: запрос к историческому снапшоту, предшествующему появлению конкретного equality или position delete file, должен видеть данные без применения этого файла - этот урок покажет, как Iceberg обеспечивает эту гарантию через тот же механизм sequence number, что был разобран в разделе про анатомию delete file.
-
Урок 10 («Table Maintenance») детально разбирает
rewrite_data_filesиrewrite_position_delete_files- процедуры, многократно анонсированные в этом уроке (модель Read Penalty, Кейс 3, производственный кейс) как единственный способ сдержать рост стоимости чтения MoR-таблицы; этот урок покажет точный синтаксис, параметры и стратегии планирования таких процедур. -
Урок 13 («Выбор формата») вернётся к сравнению CoW/MoR Iceberg с аналогичными механизмами Delta Lake и Apache Hudi (например, Hudi Copy-on-Write/Merge-on-Read tables или Delta Lake Deletion Vectors) - материал этого урока станет прямой базой для понимания, что описанный здесь компромисс не уникален для Iceberg, а является общим паттерном современных table format'ов.
Домашнее задание¶
-
Воспроизведите Кейс 1 практического блока на собственном self-hosted стенде (PostgreSQL + MinIO + Spark), но выполните
DELETEпо условию, точно совпадающему с ключом партиционирования (event_tsв границах одного дня), а не поstatus/customer_id. Сравнитеrecord_countитогового position delete file с результатом из урока и объясните разницу через partition-level pruning (четвёртый урок модуля). -
Используя
position_delete_merge_costиequality_delete_merge_costиз раздела про формальную оценку стоимости merge, постройте таблицу (текстовую или через matplotlib) зависимости расчётной стоимости merge от числа применимых delete-файлов (1, 5, 10, 25, 50) для фиксированного data-файла на 50 000 строк, отдельно для position и equality стратегий. Сделайте вывод о том, при каком числе delete-файлов равенство стоимостей становится принципиально невозможным независимо от паттерна нагрузки. -
Повторите Кейс 3, но вместо 25 точечных
DELETEвыполните пять операцийMERGE INTO(синтаксис будет детально разобран в восьмом уроке модуля, но базовая формаMERGE INTO ... WHEN MATCHED THEN UPDATEуже доступна), каждая из которых обновляет 200 случайных строк той же партиции. Сравните прирост числа delete-файлов и времяSELECTс результатами Кейса 3 и объясните разницу через число затронутых строк за одну операцию. -
Изучите (экспериментально или по документации Apache Iceberg) системную таблицу
.all_entriesв сравнении с.entries, использованной в Кейсе 3. Объясните письменно, почему.entriesпоказывает только записи манифестов текущего снапшота, тогда как.all_entriesвключает записи из всех снапшотов истории таблицы, и в каком сценарии диагностики delete-файлов второй вариант необходим. -
На основе производственного кейса урока спроектируйте (без необходимости реализации) конкретную метрику и пороговое значение для алертинга, которые позволили бы обнаружить деградацию из кейса на неделе 2-3, а не на неделе 6. Обоснуйте письменно выбор источника метрики между
.delete_files(общее число файлов) и.entries(среднее число применимых delete-файлов на data-файл). -
Создайте таблицу с
write.delete.mode = 'merge-on-read'иwrite.delete.format.default = 'avro'(вместоparquetпо умолчанию). ВыполнитеDELETE, аналогичный Кейсу 1, и через.delete_filesсравнитеfile_size_in_bytesрезультата с измеренным в уроке значением дляparquet. Письменно объясните полученную разницу через различие между построчным (Avro) и колоночным (Parquet) форматом для файла такого малого размера. -
Спроектируйте (без необходимости реализации) конфигурацию
write.update.mode/write.delete.mode/write.merge.modeдля гипотетической таблицыpayments, которая одновременно: (а) принимает CDC-поток от Debezium с частотой сотни событий в секунду, (б) читается аналитическим дашбордом каждые 30 секунд. Обоснуйте письменно, почему чисто табличной конфигурации в этом случае может быть недостаточно и какая дополнительная эксплуатационная мера (помимо выбора режима) обязательно потребуется, опираясь на материал производственного кейса. -
Напишите unit-тест (на PySpark, можно через
pytestс локальным Iceberg-каталогом), который проверяет инвариант: после любогоDELETEна MoR-таблице суммаsummary['added-position-deletes']нового снапшота равна числу строк, реально удовлетворивших условиюWHERE(проверяемому отдельнымCOUNT(*)до удаления). Это - практическая проверка корректности модели, разобранной в Кейсе 1. -
Сравните измеренное время выполнения
SELECT(как в Кейсе 3) для трёх вариантов одной и той же партиции с 25 накопленными position delete files: (а) как в Кейсе 3, без изменений; (б) после принудительногоCACHE TABLEпартиции перед измерением; (в) при увеличенииspark.sql.files.maxPartitionBytesдо значения, превышающего размер партиции. Объясните полученное ранжирование, опираясь на то, какая часть стоимости Read Penalty относится к самому merge, а какая - к обычному I/O чтения data-файла. -
Используя дерево решений из раздела про конфигурацию, письменно (1-2 абзаца) разберите следующий пограничный случай: таблица читается батч-джобом раз в час (не read-heavy в смысле BI-дашбордов), но изменяется потоковым CDC-источником непрерывно. Объясните, почему ответ дерева решений («часто» на
Q1, ведёт к MoR) в этом случае всё ещё корректен, несмотря на относительно редкое чтение, опираясь на то, что Read Penalty платится не за каждое обращение пользователя, а за каждое сканирование таблицы движком - и единственный батч-джоб раз в час всё равно обязан пройти Этап 3.1 для каждой партиции с накопленными delete-файлами.
Полная картина: жизненный цикл одной Merge-on-Read операции до компакции¶
Завершая урок, соберём весь материал в единую диаграмму - от поступления изменения до момента, когда накопленные delete-файлы наконец материализуются через компакцию (детали которой - тема десятого урока).
Эта диаграмма - зеркальное отражение капстоун-диаграммы прошлого урока: там путь заканчивался состоянием CLEAN сразу после одного commit'а, а здесь между commit'ом и состоянием CLEAN появляется целая дополнительная петля (READ -.-> ACCUM), отражающая главное отличие MoR - стоимость, которая не исчезает после записи, а накапливается и применяется заново при каждом обращении к таблице, пока её не разорвёт плановая компакция. Каждый блок этой диаграммы был детально разобран в отдельном разделе урока: ветвление SRC - в разделе про blind write equality deletes; ACCUM/READ - в моделях Read Penalty и формальной оценки стоимости merge; COMPACT/CLEAN - анонсированы как тема десятого урока модуля.
Итоги¶
Merge-on-Read - не бесплатная альтернатива Copy-on-Write, а перенос стоимости с записи на чтение. Прошлый урок показал цену немедленного переписывания файлов; этот урок показал симметричную цену - стоимость, которая не платится один раз, а накапливается и применяется заново при каждом SELECT, пока её не сбросит компакция.
Delete file - новый тип физического файла в Iceberg v2, существующий в двух формах. Position delete хранит точные координаты (file_path, pos) удаляемой строки; equality delete хранит значения ключевых колонок (equality_ids) без знания позиции. Манифесты v2 разделяют записи о data- и delete-файлах через поле content, никогда не смешивая их в одном манифесте.
Sequence number - единственный механизм, делающий применение delete-файлов корректным без блокировок. Delete-файл применяется к data-файлу, если его sequence_number больше или равен sequence_number данных в этом data-файле - правило, разобранное во втором уроке модуля и формально объясняющее, почему параллельные writer'ы не нуждаются в координации.
Position deletes требуют scan перед записью, но платят только за число изменённых строк. Spark SQL DELETE/UPDATE/MERGE всегда генерируют именно position deletes, потому что собственный scan движка делает позиции бесплатно доступными - это было прямо подтверждено в Кейсе 1 практического блока.
Equality deletes устраняют scan полностью за счёт blind write, но создают fan-out join при чтении. Это - единственный практичный механизм для потоковых CDC-writer'ов без доступа к позициям (Кейс 2), но каждый накопленный equality delete file требует join со всеми применимыми data-файлами партиции - и эта стоимость на три порядка выше аналогичного числа position deletes (раздел про формальную оценку стоимости merge).
Read Penalty растёт со временем, а не фиксируется в момент записи. Это - принципиальное отличие от разовой Write Amplification Copy-on-Write: каждая новая MoR-операция добавляет ещё один потенциально применимый delete-файл к уже существующим data-файлам, и стоимость каждого следующего чтения постепенно увеличивается, что было измерено в Кейсе 3 и подтверждено производственным кейсом урока.
Регулярная компакция - обязательное условие эксплуатации MoR-таблицы, а не опциональная оптимизация. Производственный кейс показал, что отсутствие операционного контроля за maintenance-job (а не сам факт выбора MoR) привело к деградации чтения в 25 раз за шесть недель - rewrite_data_files/rewrite_position_delete_files (десятый урок модуля) - не дополнение к архитектуре MoR-таблицы, а её неотъемлемая часть.
Выбор между CoW и MoR, а также между position и equality внутри MoR, определяется профилем нагрузки, а не предпочтением инженера. Дерево решений и сравнительная матрица этого урока формализуют то же правило, что было сформулировано в прошлом уроке: решение должно опираться на измерение частоты изменений и источника этих изменений, а не на интуитивный выбор «более современного» или «более простого» режима.
Краткий глоссарий терминов урока¶
-
Merge-on-Read (MoR) - стратегия изменения данных, при которой логическое изменение строки реализуется через запись отдельного delete-файла рядом с неизменным data-файлом, а согласование состояния происходит во время чтения.
-
Delete file - физический файл (обычно Parquet), хранящий инструкцию о том, какие строки существующих data-файлов считаются удалёнными; не содержит «полезных» данных пользователя в привычном смысле.
-
Position Delete - тип delete-файла (
content = POSITION_DELETES), хранящий точные координаты удаляемой строки в виде пары(file_path, pos). -
Equality Delete - тип delete-файла (
content = EQUALITY_DELETES), хранящий значения одной или нескольких колонок (equality_ids), идентифицирующих удаляемые строки по условию, а не по позиции. -
Blind Write - запись delete-файла без какого-либо обращения к существующим данным таблицы; характерное свойство equality deletes, делающее их пригодными для потоковых CDC-writer'ов.
-
Read Penalty - дополнительная вычислительная стоимость, которую каждый
SELECTплатит за объединение data-файлов с применимыми delete-файлами; в отличие от Write Amplification, накапливается со временем без компакции. -
equality_ids- список стабильных field id колонок, образующих ключ удаления equality delete file; ссылается на ту же систему идентификаторов, что разбиралась в уроке про partition evolution. -
Fan-out join (equality deletes) - ситуация, при которой один equality delete file должен быть применён («join'нут») к каждому применимому data-файлу партиции, поскольку delete-файл не содержит информации о конкретном местоположении совпадающих строк.
-
sequence_number/file_sequence_number- пара полей записи манифеста v2; первое определяет применимость delete-файла к data-файлу, второе отражает момент физической записи файла (обычно совпадают, расходятся приcherry-pick). -
Manifest content (
DATA/DELETES) - поле манифест-файла v2, фиксирующее, что весь манифест содержит записи только одного типа - либо только data-файлы, либо только delete-файлы. -
.delete_files- системная метаданных-таблица Iceberg Spark extensions, специфичная для v2, показывающая все активные delete-файлы таблицы с колонкамиcontent,equality_ids,record_countи другими. -
.entries- системная таблица, показывающая записи манифестов текущего снапшота, включая полеdata_file.content, используемое для подсчёта числа применимых delete-файлов на data-файл. -
write.delete.format.default- свойство таблицы, определяющее формат файла для delete-файлов (по умолчаниюparquetв актуальных версиях Apache Iceberg при записи через Spark). -
PositionDeleteIndex (концептуальная модель) - структура, которую executor строит в памяти при чтении применимых position delete files для эффективной построчной фильтрации data-файла.
-
LSM-дерево (Log-Structured Merge-tree) как аналогия - структура данных (Cassandra, RocksDB), в которой запись никогда не модифицирует существующие сегменты, а изменения и tombstone'ы добавляются как новые append'ы, объединяемые при чтении - прямая концептуальная параллель с архитектурой Merge-on-Read.
-
Tombstone (аналогия из LSM-дерева) - запись, помечающая логическое удаление без физического удаления данных; delete file Iceberg - табличный аналог tombstone-записи.
-
Компакция (предварительное упоминание, детали в десятом уроке) - процесс материализации накопленных delete-файлов обратно в чистые data-файлы (
rewrite_data_files/rewrite_position_delete_files), являющийся прямым аналогом LSM-compaction и обязательным условием долгосрочной эксплуатации MoR-таблицы. -
Read Amplification (формальная оценка стоимости merge) - количественная мера дополнительной работы, выполняемой при чтении из-за применимых delete-файлов; растёт линейно с числом удалённых строк для position deletes и мультипликативно с числом применимых delete-файлов для equality deletes.
-
Антипаттерн «equality deletes без компакции на широкой партиции» - ситуация, при которой накопление equality delete files на партицию с большим числом data-файлов создаёт fan-out join, чья стоимость растёт на порядки быстрее, чем аналогичное накопление position deletes - детально измеренная в производственном кейсе урока.
-
Read-heavy/Write-heavy профиль нагрузки (критерий выбора режима) - ключевая характеристика таблицы, определяющая выбор между CoW (write-heavy цена допустима, read должен быть бескомпромиссным) и MoR (write должен быть почти бесплатным, read может себе позволить дополнительную стоимость при условии регулярной компакции).
-
GenericAppenderFactory/RowDelta(низкоуровневый Java API) - классы Iceberg Core API, используемые потоковыми коннекторами (Flink) и продвинутыми инженерами для прямой записи delete-файлов и атомарного commit'а изменений в обход стандартного SQL-планировщика row-level операций; были использованы для иллюстрации blind write в Кейсе 2. -
Дерево решений CoW vs MoR (Position) vs MoR (Equality) - последовательность из трёх вопросов (частота изменений, источник изменений, доступность позиций), формализующая инженерный выбор режима таблицы, разобранный в разделе про конфигурацию этого урока.
-
Read Penalty vs Write Amplification (симметрия двух уроков) - парные метрики стоимости двух стратегий: Write Amplification (прошлый урок) фиксируется в момент commit'а CoW-операции и не меняется постфактум; Read Penalty (этот урок) накапливается со временем и сбрасывается только плановой компакцией.