MERGE INTO: upsert паттерн для CDC - структура запроса и настройки
MERGE INTO в Apache Iceberg: синтаксис WHEN MATCHED/WHEN NOT MATCHED, физический план row-level операции под CoW и MoR, настройка write.merge.mode и write.distribution-mode, дедупликация CDC-потока и идемпотентный upsert-пайплайн на PySpark.
Зачем нужен MERGE INTO: от двух стратегий записи к production CDC-паттерну¶
Два прошлых урока разобрали Copy-on-Write и Merge-on-Read - две физические стратегии того, как Iceberg реализует построчное изменение данных внутри одного типа операции (UPDATE или DELETE по отдельности). Но реальные production-пайплайны почти никогда не сталкиваются с одной изолированной операцией - они синхронизируют целевую Iceberg-таблицу с источником, который непрерывно поставляет смесь вставок, обновлений и удалений одновременно, и делает это в виде потока изменений, а не разовой выгрузки.
Эта задача называется CDC-репликация (Change Data Capture), и именно для неё спроектирован MERGE INTO - SQL-оператор, который за один атомарный коммит способен одновременно вставить новые строки, обновить изменившиеся и удалить исчезнувшие. Этот урок - кульминация всего модуля: он связывает воедино дерево метаданных (второй урок), Scan Planning (третий урок), Hidden Partitioning (четвёртый урок) и обе стратегии построчного изменения (шестой и седьмой уроки) в единый production-паттерн, ради которого компании чаще всего и мигрируют с «голого» Parquet на Apache Iceberg.
Почему INSERT OVERWRITE и APPEND не решают задачу синхронизации¶
Прежде чем переходить к MERGE INTO, стоит явно разобрать, почему два более простых, уже знакомых инструмента - INSERT OVERWRITE и APPEND - не годятся для CDC, несмотря на то, что технически оба способны записать данные в Iceberg-таблицу.
INSERT OVERWRITE партиции решает задачу «полностью пересобрать раздел таблицы», но для этого пайплайн обязан заранее знать полное актуальное состояние раздела - то есть фактически обязан сам выполнить всю работу по слиянию старых и новых данных до записи, средствами, внешними по отношению к Iceberg. Если источник поставляет только дельту (как делает любая CDC-система), INSERT OVERWRITE просто не имеет входных данных для этой операции - ему нужен полный снимок, а не дельта.
APPEND решает противоположную проблему не лучше: новые события дописываются как новые строки, но старая версия той же логической записи никуда не исчезает. Результат - не таблица с актуальным состоянием, а журнал всех версий каждой строки, в котором обычный SELECT без дополнительной логики дедупликации увидит несколько «активных» версий одного и того же бизнес-ключа. Это разрушает базовое свойство, которое аналитики и BI-инструменты ожидают от таблицы измерений или таблицы фактов с уникальным ключом - ровно одну актуальную строку на одну сущность.
MERGE INTO - единственный из трёх путей, который принимает на вход именно то, что реально поставляет CDC-источник: дельту, а не полный снимок, - и при этом гарантирует ровно одну актуальную строку на бизнес-ключ после применения. Это и есть Upsert (от слов Update + Insert): если строка с данным ключом уже существует в целевой таблице - она обновляется; если не существует - вставляется как новая. К Upsert почти всегда добавляется и обработка удалений, поэтому в production-литературе паттерн иногда уточняют как «Upsert + Delete», хотя в синтаксисе MERGE INTO все три действия выражаются одним и тем же оператором.
CDC и событийная модель: что на самом деле приходит на вход MERGE INTO¶
Чтобы написать корректный MERGE INTO, нужно понимать формат и гарантии источника, который поставляет дельту. В большинстве production-архитектур таким источником служит связка Debezium + Kafka (или managed-аналог), реализующая CDC через чтение журнала транзакций исходной СУБД.
Как Debezium формирует событие изменения¶
Debezium подключается к write-ahead log PostgreSQL (через pgoutput/wal2json плагин логической репликации) или к binlog MySQL и читает поток уже зафиксированных транзакций - без единого SELECT к самой таблице-источнику, то есть без дополнительной нагрузки на OLTP-базу сверх той, что она уже несёт ради собственной репликации. Каждое изменение строки превращается в событие со стандартизированной структурой:
{
"op": "u",
"ts_ms": 1718880000123,
"before": {
"order_id": 1001, "customer_id": 42,
"status": "created", "amount": "199.90"
},
"after": {
"order_id": 1001, "customer_id": 42,
"status": "paid", "amount": "199.90"
},
"source": { "lsn": 24831920, "table": "orders" }
}
Поле op принимает значения c (create/insert), u (update), d (delete) и r (read - первоначальный snapshot существующих данных при первом запуске коннектора). before/after содержат полный образ строки до и после изменения (для d присутствует только before, для c - только after). Поле lsn (log sequence number, аналог sequence_number Iceberg, но на уровне WAL источника) задаёт строгий порядок событий внутри одной таблицы-источника.
Почему обработка должна быть атомарной для целого микробатча¶
Spark Structured Streaming читает поток событий из Kafka микробатчами: каждый триггер обрабатывает не одно событие, а пачку, которая может содержать вставки, обновления и удаления для разных строк одновременно. Если применять их по отдельности - например, сначала все INSERT, затем все UPDATE, затем все DELETE тремя разными командами - между этими тремя командами окажутся промежуточные, частично применённые состояния таблицы, видимые любому параллельному читателю в этот момент. Второй урок модуля показал, что снимок Iceberg-таблицы атомарен только в границах одного коммита - а значит, единственный способ избежать промежуточных несогласованных состояний - это применить весь микробатч целиком за один коммит. Именно так спроектирован MERGE INTO: он принимает источник целиком (он может содержать любую смесь типов событий) и производит ровно один новый снимок.
Промежуточная стейджинг-таблица orders_cdc на диаграмме - не обязательный, но широко используемый production-паттерн: вместо того чтобы строить MERGE INTO напрямую над потоковым DataFrame (что технически возможно через foreachBatch, но усложняет отладку и переигровку), микробатч сначала материализуется как обычная таблица (Iceberg или временная), а затем MERGE INTO выполняется уже над ней - это даёт возможность инспектировать, переигровать и дедуплицировать дельту до применения, что станет центральной темой раздела про каверзные сценарии CDC далее в этом уроке.
Анатомия и синтаксис команды MERGE INTO в Spark SQL¶
Синтаксис MERGE INTO, который реализуют Iceberg Spark extensions, следует стандарту ANSI SQL:2003 (тому же, который исторически появился в Oracle и затем был принят большинством аналитических СУБД) - это намеренное архитектурное решение: инженер, уже знакомый с MERGE в Oracle, SQL Server или Snowflake, может перенести этот навык на Iceberg почти без изменений.
Базовая структура запроса¶
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
Разберём эту конструкцию по частям:
-
MERGE INTO ... t- целевая (target) таблица, та самая Iceberg-таблица, чьё состояние должно стать согласованным с источником. Алиасtиспользуется во всех последующих условиях для однозначной ссылки на колонки целевой таблицы. -
USING orders_cdc s- источник (source) дельты. Это не обязательно физическая таблица:USINGпринимает любой именованный реляционный объект - представление (VIEW), под запрос (SELECT ... ) s) или CTE. На практике источник почти всегда оказывается результатом дополнительной обработки сырых CDC-событий (дедупликация, фильтрация), а не самой сырой стейджинг-таблицей - это станет центральной темой раздела про обработку каверзных сценариев. -
ON t.order_id = s.order_id- условие сопоставления, определяющее бизнес-ключ, по которому строки целевой таблицы и источника считаются «одной и той же сущностью». От корректности и избирательности этого условия зависит абсолютно всё остальное в запросе - если ключ не уникален в источнике или в целевой таблице, результатMERGEстановится недетерминированным (этот сценарий разберёт отдельный раздел про дубликаты). -
WHEN MATCHED THEN UPDATE SET *- действие для строк, у которых нашлась пара поON-условию: обновить целевую строку значениями из source. Звёздочка - сокращённая форма, требующая идентичного набора имён колонок в target и source (без посторонних технических полей вродеop, о чём отдельно ниже). -
WHEN NOT MATCHED THEN INSERT *- действие для строк source, для которых пары в target не нашлось: вставить как новую строку.
Любая корректная Iceberg-таблица, участвующая в MERGE INTO, должна быть таблицей формата v2 (как и для построчных UPDATE/DELETE, разобранных в прошлых двух уроках) - первая версия спецификации Iceberg не поддерживает row-level операции вообще, потому что не имеет delete-файлов как концепции.
WHEN MATCHED THEN UPDATE SET: явный список колонок против звёздочки¶
UPDATE SET * удобен, но скрывает важное ограничение: он работает только тогда, когда схема source по именам колонок в точности соответствует тому подмножеству колонок target, которое требуется обновить. В реальных CDC-пайплайнах source почти всегда содержит технические колонки, которых нет в target (op, _kafka_offset, ts_ms, lsn) - и звёздочка в этом случае либо упадёт с ошибкой схемы, либо (в зависимости от версии движка и настроек) попытается записать лишние колонки, которых не существует в target, и тоже завершится ошибкой.
Production-практика - использовать явный список присвоений, который одновременно служит документацией того, какие поля реально подлежат изменению:
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET
t.status = s.status,
t.amount = s.amount,
t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, updated_at)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.updated_at)
Явный список INSERT (...) VALUES (...) решает ту же проблему для ветки вставки: технические колонки источника (op, lsn) просто не упоминаются и не попадают в целевую таблицу. Это - не избыточная многословность, а защита от самого частого практического сбоя при подключении нового CDC-источника: несовпадения схем source и target, которое звёздочка превращает в ошибку времени выполнения вместо явной, читаемой спецификации в самом тексте запроса.
WHEN MATCHED THEN DELETE: обработка удалений из источника¶
CDC-событие с op = 'd' означает, что строка была физически удалена в источнике - и целевая Iceberg-таблица должна отразить это тем же DELETE, что разбирался в шестом и седьмом уроках, но теперь - как одна из веток единого MERGE INTO, а не отдельная команда:
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc s
ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'd' THEN DELETE
WHEN MATCHED AND s.op != 'd' THEN UPDATE SET
t.status = s.status,
t.amount = s.amount,
t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, updated_at)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.updated_at)
Альтернатива физическому DELETE - soft delete: вместо удаления строки целевая таблица хранит флаг is_deleted = true, а сама строка остаётся видимой для аналитики истории (что часто требуется для аудита или построения SCD Type 2 - тема, к которой курс возвращается в модуле про моделирование данных). Выбор между физическим и логическим удалением - архитектурное решение, не зависящее от MERGE INTO как механизма; оператор одинаково легко выражает оба варианта.
WHEN NOT MATCHED THEN INSERT: вставка новых записей¶
Ветка INSERT обрабатывает строки source, для которых не нашлось пары в target по ON-условию - то есть новые сущности, которых раньше не существовало в целевой таблице. Как и UPDATE, она может принимать условие (WHEN NOT MATCHED AND s.op = 'c' THEN INSERT ...), что используется реже, поскольку для отсутствующей в target строки операция логически всегда сводится к вставке независимо от исходного op (даже op = 'u' для несуществующей строки означает «вставить её as is» - ситуация, которая возникает, например, при первом запуске пайплайна на таблице, где исходный CDC-снапшот ещё не был догружен).
Множественные предикаты и порядок WHEN-веток¶
Iceberg позволяет указать произвольное число WHEN MATCHED/WHEN NOT MATCHED веток в одном запросе, и они оцениваются строго по порядку сверху вниз, как ветки CASE WHEN - применяется первая, чьё условие истинно, и обработка строки на этом завершается. Это даёт возможность встраивать защитные бизнес-условия прямо в структуру запроса:
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc s
ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'd' THEN DELETE
WHEN MATCHED AND s.updated_at > t.updated_at THEN UPDATE SET
t.status = s.status,
t.amount = s.amount,
t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, updated_at)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.updated_at)
Здесь вторая ветка UPDATE сработает только тогда, когда событие source действительно новее, чем текущее состояние target (s.updated_at > t.updated_at). Критически важная деталь, часто упускаемая при первом знакомстве с MERGE INTO: если строка совпала по ON-условию (то есть формально является MATCHED), но ни одна из перечисленных WHEN MATCHED-веток не выполнила своё условие - целевая строка остаётся полностью нетронутой, без какой-либо ошибки или неявного действия. Это и есть механизм, который делает защиту от устаревших событий (раздел про late-arriving data) корректной: опоздавшее событие просто молча игнорируется для уже актуальной строки, вместо того чтобы откатить её состояние назад или прервать выполнение всего батча ошибкой.
Отдельная техническая оговорка про стандарт ANSI: начиная со Spark 3.4, парсер Catalyst поддерживает третью форму ветки - WHEN NOT MATCHED BY SOURCE, описывающую строки target, отсутствующие в source (применяется, например, для полной сверки с периодическим full snapshot, а не с CDC-дельтой). Поддержка этой ветки конкретным движком зависит от того, реализует ли используемая версия Iceberg-Spark runtime соответствующий интерфейс SupportsRowLevelOperations для этого случая - перед использованием в production стоит проверить changelog установленной версии iceberg-spark-runtime. Для классического CDC-паттерна этого урока (источник - дельта, а не полный снимок) эта ветка обычно не нужна: «удаление, отсутствующее в источнике» в CDC-модели выражается явным событием op = 'd', а не отсутствием строки в source.
Физический план выполнения MERGE INTO под капотом Spark¶
Текст SQL-запроса одинаков независимо от того, как Iceberg физически выполнит каждую его ветку - но реальная нагрузка на кластер определяется именно физическим планом, который строит планировщик row-level операций. Этот раздел разбирает, во что Catalyst превращает MERGE INTO, прежде чем передать план executor'ам.
Row-level operations API: почему MERGE INTO - не эмуляция, а нативная операция каталога¶
Начиная с DataSource V2 API, Spark определяет интерфейс SupportsRowLevelOperations, который источник данных (в данном случае - Iceberg-каталог) может реализовать, чтобы взять на себя планирование UPDATE/DELETE/MERGE самостоятельно, вместо того чтобы Spark эмулировал их через generic «прочитать всё - перезаписать всё». Iceberg реализует этот интерфейс, и именно поэтому шестой и седьмой уроки могли показать два принципиально разных физических результата (переписанный файл при CoW, маленький delta-файл при MoR) для синтаксически одинаковых команд - выбор стратегии полностью находится внутри Iceberg-каталога, а не в общей логике Spark SQL.
MERGE INTO использует тот же интерфейс, но с дополнительной сложностью: в отличие от UPDATE/DELETE, у которых ровно один источник правды (сама целевая таблица плюс условие WHERE), MERGE INTO должен сопоставить две таблицы - и именно это сопоставление реализуется как join.
Join-стратегия: как Spark сопоставляет target и source¶
Целевая таблица перед join'ом дополняется служебными метаданными-колонками _file и _pos (то же понятие физической позиции строки, что лежит в основе position delete file из седьмого урока) - они присутствуют в плане выполнения, но не видны в результате SELECT * FROM target. Это даёт планировщику возможность точно идентифицировать каждую строку target по её физическому местоположению, не дожидаясь отдельного шага поиска позиций, как это было бы для plain DELETE с произвольным условием.
Логический тип join, который должен быть построен, зависит от того, какие именно WHEN-ветки присутствуют в запросе - и это решение определяется самой природой каждой ветки, а не настройкой конфигурации:
Логика каждой ветви этого дерева напрямую следует из того, какие строки обязаны «выжить» в результате join, чтобы дальнейшая обработка была корректной:
-
Если в запросе только
WHEN MATCHED-ветки (UPDATE/DELETE, безINSERT) - план должен сохранить все строки target (а не только совпавшие), потому что, как будет показано в следующем подразделе, при CoW необходимо знать содержимое каждой строки физически затронутого файла, включая несовпавшие. Это требуетRIGHT OUTER JOIN, в котором target - сохраняемая сторона. -
Если присутствует и
WHEN MATCHED, иWHEN NOT MATCHED- план должен дополнительно опознать строки source, не нашедшие пары в target (кандидаты наINSERT). Это требуетFULL OUTER JOIN: сохраняются и непарные строки target, и непарные строки source. -
Если в запросе только
WHEN NOT MATCHED(чистый дедуплицирующийINSERT, без единойUPDATE/DELETE-ветки) - единственное, что нужно найти - строки source без пары в target; сами строки target в обработке не участвуют вообще. Это самый дешёвый случай, требующий лишьLEFT ANTI JOIN.
Уточнение: чем join для CoW отличается от join для MoR при одинаковом тексте запроса¶
Дерево выше показывает логический join, продиктованный структурой SQL-запроса - но шестой и седьмой уроки уже показали, что именно физически читается и записывается, может кардинально отличаться в зависимости от write.merge.mode. Это отличие стоит сделать явным именно для MERGE INTO, потому что на первый взгляд кажется нелогичным: SQL-текст один и тот же, План выполнения (join) выше один и тот же логически - а физическая работа executor'ов разная.
При Copy-on-Write физический результат join'а должен содержать все строки затронутого файла, потому что новый data-файл, который заменит старый, обязан содержать корректную версию каждой строки - совпавшие обновляются, несовпавшие копируются без изменений байт в байт. Это прямое продолжение Read-Modify-Write модели из шестого урока, но теперь применённое не к одной операции, а к трём действиям (UPDATE/DELETE/INSERT) одновременно, дающим один новый файл на выходе.
При Merge-on-Read логический join остаётся тем же FULL OUTER JOIN (план всё ещё должен опознать оба типа «сирот»), но физическая запись принципиально иная: несовпавшие строки затронутого файла не нужно копировать никуда - они просто остаются на своих местах в исходном файле, нетронутые. Executor'у достаточно записать координаты совпавших строк в position delete file (ровно тот механизм, что разобран в седьмом уроке) и новые версии этих строк - вместе с вставленными - в отдельный, обычно очень маленький новый data-файл. Именно эта экономия физической работы - прямое объяснение того, почему MERGE INTO в режиме MoR коммитится быстрее и почему производственный кейс этого урока (далее) выберет именно этот режим для высокочастотного CDC-потока.
Проблема перекоса данных (Data Skew) в MERGE INTO¶
Join, лежащий в основе MERGE INTO, подвержен той же проблеме, что и любой другой join в Spark: если распределение значений ключа ON-условия неравномерно, часть executor-задач получает кратно больше работы, чем остальные, и становится «стрэглерами», определяющими общее время выполнения всей операции.
Для CDC-нагрузок перекос особенно вероятен по двум причинам. Во-первых, бизнес-данные сами по себе часто скошены - например, небольшая доля customer_id (крупные корпоративные клиенты) генерирует значительно больше заказов, чем типичный розничный покупатель, и именно эти строки чаще встречаются в обновлениях. Во-вторых, если целевая таблица партиционирована по одному измерению (event_ts, как в примерах шестого и седьмого уроков), а join выполняется по другому ключу (customer_id/order_id), пруннинг по партициям (четвёртый урок) не способен сократить перекос - он работает на уровне файлов, а не на уровне распределения значений ключа join внутри executor-задач.
Симптом перекоса хорошо виден в Spark UI: в стадии (stage), соответствующей join'у MERGE INTO, гистограмма длительности задач показывает один или несколько резко выделяющихся «хвостов», тогда как большинство задач завершается быстро. Механизм автоматического выявления и обработки именно такого перекоса в join - Adaptive Query Execution - разбирался отдельно в модуле про партиционирование; для MERGE INTO критично, чтобы spark.sql.adaptive.skewJoin.enabled оставался включённым (по умолчанию - true начиная со Spark 3), поскольку без него единственным средством борьбы с перекосом остаётся ручная техника соления ключа (salting), требующая изменения самого запроса. Дополнительный, специфичный для Iceberg рычаг - настройка write.distribution-mode, разбираемая в следующем разделе: принудительный shuffle перед записью способен заметно смягчить последствия перекоса на стороне записи, даже если сам join всё ещё выполняется по скошенному ключу.
Сравнительная таблица: одна и та же команда MERGE INTO под CoW и MoR¶
Прежде чем переходить к настройке через TBLPROPERTIES, сведём в одну таблицу всё, что отличает физическое исполнение одного и того же текста запроса MERGE INTO в зависимости от выбранного режима - своего рода краткая выжимка всего раздела про физический план, к которой удобно возвращаться при проектировании новой таблицы:
| Параметр | Copy-on-Write | Merge-on-Read |
|---|---|---|
| Что физически читается | Весь затронутый data-файл целиком | Только строки, нужные для join (сам файл сканируется, но не копируется) |
| Что физически пишется | Новый data-файл, дублирующий все строки старого (изменённые + неизменные + вставленные) | Маленький data-файл (новые версии + inserts) и position delete file (координаты старых версий) |
| Несовпавшие строки старого файла | Копируются байт в байт в новый файл | Не читаются повторно, остаются на месте без изменений |
| Стоимость записи на одно изменение | Пропорциональна размеру файла, не строки (Write Amplification) | Пропорциональна числу изменённых строк, почти не зависит от размера файла |
| Стоимость чтения после MERGE | Не изменяется - файл уже консолидирован | Растёт с числом накопленных delete-файлов (Read Penalty, седьмой урок) |
| Логический join (этот раздел) | Тот же RIGHT OUTER/FULL OUTER/LEFT ANTI, что и для MoR |
Тот же RIGHT OUTER/FULL OUTER/LEFT ANTI, что и для CoW |
| Типичный профиль нагрузки | Редкие batch-обновления, read-heavy таблицы (BI, Gold-слой) | Частые потоковые CDC-микробатчи, write-heavy профиль |
| Обязательная maintenance-операция | Не требуется для самого MERGE (но полезна для общего числа файлов) | Регулярная компакция обязательна (десятый урок модуля) |
Главный вывод этой таблицы - логический план (строка про join) у CoW и MoR совпадает полностью; вся разница сосредоточена в строках, описывающих физическое чтение и запись. Это - важное концептуальное уточнение к материалу предыдущего раздела: write.merge.mode не меняет что вычисляется, а меняет только как много байт физически перемещается ради того же самого логического результата.
Тонкая настройка производительности через TBLPROPERTIES¶
Помимо логики самого запроса, Iceberg даёт три независимых рычага конфигурации таблицы, которые напрямую определяют, насколько дорогим окажется регулярный MERGE INTO под CDC-нагрузкой. Все три устанавливаются один раз на уровне таблицы и не требуют изменения самого SQL-запроса.
write.merge.mode: тот же переключатель CoW/MoR, но именно для MERGE¶
Шестой и седьмой уроки представили три независимые table properties - write.update.mode, write.delete.mode, write.merge.mode - каждая из которых настраивает свой тип row-level операции независимо от остальных. Это разделение не теоретическое: реальная таблица может, например, использовать merge-on-read для MERGE INTO (частый CDC-поток) и одновременно copy-on-write для редких ручных UPDATE (например, разовая корректировка данных аналитиком) - конфигурация каждой операции отражает именно её собственный профиль нагрузки.
CREATE TABLE lakehouse.analytics.orders (
order_id BIGINT,
customer_id BIGINT,
status STRING,
amount DECIMAL(10, 2),
updated_at TIMESTAMP,
event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(event_ts))
TBLPROPERTIES (
'format-version' = '2',
'write.merge.mode' = 'merge-on-read'
);
Если свойство write.merge.mode не указано явно, Iceberg использует значение по умолчанию, совпадающее со значением по умолчанию для write.update.mode и write.delete.mode (в актуальных версиях - copy-on-write) - то есть поведение по умолчанию консервативно и предсказуемо, но для регулярного высокочастотного CDC-потока почти всегда требует явного переопределения в merge-on-read, опираясь на тот же инженерный анализ, что был разобран в шестом и седьмом уроках.
write.distribution-mode: борьба с проблемой мелких файлов¶
Третий урок модуля показал, что разрастание числа манифестов и манифест-листов - не бесплатная операция: каждый дополнительный data-файл означает дополнительную запись в манифесте, а значит - дополнительную работу при Scan Planning любого будущего запроса. MERGE INTO, выполняемый часто (как и положено CDC-пайплайну), без специальной настройки имеет тенденцию создавать ровно такую проблему: каждый микробатч порождает собственный набор файлов, и если эти файлы недостаточно крупные и недостаточно упорядоченные по партициям, число файлов в таблице растёт быстрее, чем растут реальные данные.
Свойство write.distribution-mode управляет тем, выполняется ли shuffle перед фактической записью данных executor'ами, и если да - по какому принципу:
ALTER TABLE lakehouse.analytics.orders SET TBLPROPERTIES (
'write.distribution-mode' = 'hash'
);
-
none- данные пишутся как есть, без дополнительного shuffle: каждая executor-задача записывает в любые партиции, которые встретились в обрабатываемых ей строках. Самый дешёвый вариант с точки зрения CPU, но при широком распределении партиций внутри одного микробатча почти гарантированно создаёт несколько мелких файлов на одну и ту же партицию - каждая задача, «зацепившая» хотя бы одну строку этой партиции, создаёт собственный файл. -
hash- перед записью выполняется shuffle по хешу от колонок партиционирования: все строки одной партиции гарантированно собираются в одну (или предсказуемо малое число) executor-задачу, которая и пишет для неё единственный консолидированный файл. Это - именно тот механизм, который превращает CDC-микробатч, разбросанный по множеству партиций, в компактный набор файлов вместо десятков мелких фрагментов. -
range- данные распределяются по диапазонам значений (через сэмплирование), а не по хешу; используется реже, в основном когда таблица также задаёт явный sort order, и важно физически сохранить упорядоченность строк внутри получившихся файлов, а не только их группировку по партициям.
Цена дополнительного shuffle - вполне реальная, и для очень маленьких микробатчей (характерных для high-frequency CDC, где один триггер обрабатывает лишь несколько тысяч строк) накладные расходы на сериализацию и передачу данных между executor'ами при hash-distribution могут превысить выгоду от уменьшения числа файлов. Production-практика - измерять обе конфигурации на реальном профиле микробатчей и явно проверять рост числа манифестов и data-файлов со временем (та же системная таблица .files, что использовалась в практических блоках шестого и седьмого уроков), а не выбирать hash по умолчанию «на всякий случай».
Bloom-фильтры: ускорение фазы поиска соответствий¶
Третий урок модуля показал, что Scan Planning отсекает нерелевантные файлы по статистике min/max на уровне манифеста, а внутри уже выбранного файла Parquet может дополнительно отсекать целые row groups по тем же min/max границам, сохранённым в footer. Но для join-операции внутри MERGE INTO, где условие ON обычно представляет собой точное равенство по высококардинальному ключу (order_id, customer_id), границы min/max часто оказываются слишком грубыми - диапазон конкретного row group может быть широким, а искомое значение ключа - единственным в этом диапазоне.
Для такого паттерна (точечный point-lookup по значению внутри уже выбранного файла) Parquet поддерживает bloom-фильтры - компактные вероятностные структуры данных, отвечающие на вопрос «может ли это значение присутствовать в данном row group» с гарантированным отсутствием false negative (если фильтр говорит «нет» - значения точно нет) и небольшой вероятностью false positive (фильтр может сказать «может быть» для значения, которого на самом деле нет, что приводит к лишнему, но не ошибочному, чтению).
ALTER TABLE lakehouse.analytics.orders SET TBLPROPERTIES (
'write.parquet.bloom-filter-enabled.column.customer_id' = 'true',
'write.parquet.bloom-filter-max-bytes' = '1048576'
);
Принципиально важно понимать порядок, в котором применяются разные уровни отсечения при MERGE INTO: сначала манифест-уровневый пруннинг (третий урок) выбирает кандидат-файлы по диапазону партиций и статистике min/max самого файла; затем, уже внутри каждого выбранного файла, bloom-фильтр позволяет пропустить целые row groups, для которых фильтр гарантированно отвечает «значения ключа join точно нет» - без необходимости декодировать и разворачивать содержимое column chunk. Bloom-фильтр не заменяет пруннинг по статистике, а дополняет его на следующем, более детальном уровне детализации - и выигрыш от него заметен именно на тех таблицах, где join probe-сторона (то есть target) содержит десятки и сотни тысяч строк на один файл, а матчей по конкретному микробатчу относительно немного.
Обработка каверзных сценариев CDC: дубликаты и опоздавшие события¶
Теория предыдущих разделов предполагала идеальный источник - дельту, в которой ровно одна строка приходится на один изменившийся бизнес-ключ. Реальные CDC-потоки этому условию не подчиняются почти никогда, и MERGE INTO обладает встроенным механизмом защиты именно от этого несоответствия - но прежде чем научиться с ним работать, нужно понять, как он проявляется в виде ошибки выполнения.
Проблема дубликатов в источнике: cardinality check¶
Представим, что в один и тот же микробатч CDC-потока попали два разных события по одному и тому же order_id - например, Kafka Connect sink с гарантией at-least-once доставил одно и то же событие повторно после кратковременного сетевого сбоя, либо исходная транзакция в OLTP-базе действительно дважды изменила одну и ту же строку в пределах окна одного микробатча. Источник orders_cdc в этом случае содержит две строки с одинаковым order_id, и ON t.order_id = s.order_id сопоставит обе этих строки источника с одной и той же строкой target.
Это - логически неоднозначная ситуация: какое из двух событий должно определить итоговое содержимое строки t? Если позволить выполнить оба UPDATE друг за другом без явного порядка, результат зависит от случайного порядка обработки внутри executor-задачи - то есть становится недетерминированным, а недетерминированные изменения данных недопустимы для production-системы, претендующей на согласованность.
Iceberg защищается от этой ситуации явной проверкой - cardinality check, включённой по умолчанию (write.merge.cardinality-check.enabled = true): если планировщик обнаруживает, что одна и та же строка target должна была бы быть изменена более чем одной строкой source в рамках одного WHEN MATCHED, выполнение всего MERGE INTO прерывается исключением времени выполнения, по смыслу эквивалентным «строка целевой таблицы сопоставлена с несколькими строками источника» (точная формулировка текста исключения отличается между версиями Spark и Iceberg, но смысл проверки одинаков во всех актуальных релизах).
spark.sql("""
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc_dirty s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET
t.status = s.status, t.amount = s.amount, t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, updated_at)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.updated_at)
""")
# org.apache.spark.SparkException: ... a single row from the target table
# was matched with multiple rows of the source table - MERGE INTO операция
# прервана, изменения НЕ применены, таблица осталась в состоянии до commit'а
Важная гарантия, прямо следующая из атомарности коммита (второй урок модуля): если cardinality check сработал, вся операция MERGE INTO целиком откатывается - ни одна строка из этого микробатча не применяется, в том числе те, что были однозначны и не имели дублей. Это - намеренное архитектурное решение в пользу согласованности данных, а не частичного, потенциально некорректного прогресса: лучше не применить весь батч и дать пайплайну возможность переиграть его после исправления источника, чем применить часть изменений в неопределённом порядке.
Паттерн дедупликации «на лету»: ROW_NUMBER() в источнике¶
Решение, диктуемое самой природой проблемы, - гарантировать уникальность бизнес-ключа в source до того, как он попадёт в USING-секцию MERGE INTO. Идиоматичный способ сделать это в Spark SQL - оконная функция ROW_NUMBER(), ранжирующая дубликаты внутри партиции окна по выбранному критерию «свежести» и оставляющая только победителя:
WITH dedup_source AS (
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY order_id
ORDER BY updated_at DESC
) AS rn
FROM orders_cdc
)
MERGE INTO lakehouse.analytics.orders t
USING (SELECT * FROM dedup_source WHERE rn = 1) s
ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'd' THEN DELETE
WHEN MATCHED AND s.updated_at > t.updated_at THEN UPDATE SET
t.status = s.status, t.amount = s.amount, t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, updated_at)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.updated_at)
PARTITION BY order_id ORDER BY updated_at DESC присваивает значение 1 той версии каждого order_id, у которой updated_at максимален - то есть самой свежей версии события в этом микробатче; внешний WHERE rn = 1 отбирает только эти строки, гарантируя ровно одну строку на ключ перед тем, как source попадёт в ON-условие. Catalyst инлайнит CTE в общий логический план запроса, поэтому дополнительные накладные расходы сводятся к самой оконной функции (требующей shuffle по order_id, если источник до этого не был распределён по этому ключу) - сопоставимо по стоимости с тем же join, который и так должен произойти для самого MERGE INTO.
Принципиальное предостережение: SELECT DISTINCT не решает ту же задачу. DISTINCT устраняет полностью идентичные строки, но если два события по одному order_id отличаются хотя бы одним полем (а они почти всегда отличаются - иначе событие не было бы сгенерировано CDC-источником вообще), DISTINCT оставит обе строки как разные, и cardinality check всё равно сработает. Дедупликация по бизнес-ключу с явным критерием упорядочивания (а не по полному совпадению строк) - единственный корректный механизм для этой задачи.
Late-Arriving Data: защита от перезаписи более новой версии устаревшим событием¶
Дедупликация внутри одного микробатча решает только часть проблемы. Представим, что событие E2 (статус paid, updated_at = 10:00:05) было успешно применено предыдущим микробатчем, а следующий микробатч, по какой-то причине (повторная доставка из Kafka, ручной replay части топика для восстановления после сбоя), приносит более старое событие E3 (статус cancelled, updated_at = 10:00:03) - событие, которое внутри своего собственного микробатча уникально по ключу и проходит дедупликацию без единого предупреждения, но по отношению к уже накопленному состоянию target оказывается устаревшим.
Без guard-условия s.updated_at > t.updated_at (раздел про множественные предикаты выше) такой UPDATE молча откатил бы видимый статус заказа с paid обратно на cancelled - регрессия, которую крайне трудно диагностировать постфактум, потому что сама операция выполняется без единой ошибки. С guard-условием опоздавшее событие просто не проходит ни одну WHEN MATCHED-ветку (как было показано в разделе про порядок веток) и оставляет строку target нетронутой - именно так, как должно быть, поскольку target уже отражает более новое состояние, чем то, что несёт опоздавшее событие.
Важно подчеркнуть, что внутрибатчевая дедупликация (ROW_NUMBER()) и межбатчевая защита от устаревших событий (updated_at-guard) решают разные проблемы и обязаны применяться вместе: первая гарантирует, что source не содержит внутренних противоречий (без неё сработает cardinality check); вторая гарантирует, что даже корректный, уникальный source не откатит состояние, которое уже новее, чем он сам (без неё данные тихо регрессируют). Применение только одной из двух защит оставляет класс ошибок, который вторая защита закрывает.
Практический демо-блок: MERGE INTO от идемпотентного upsert до провокации ошибки¶
Стенд использует ту же self-hosted инфраструктуру, что и весь модуль: JDBC Catalog поверх PostgreSQL для метаданных и S3FileIO поверх MinIO для физических файлов (полная конфигурация SparkSession приведена в первом уроке модуля и здесь не повторяется). Все примеры ниже выполняются в pyspark с подключёнными Iceberg Spark extensions, и продолжают работу с той же таблицей lakehouse.analytics.orders, что и в шестом и седьмом уроках.
Кейс 0: подготовка целевой таблицы и первой версии CDC-источника¶
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.analytics.orders (
order_id BIGINT,
customer_id BIGINT,
status STRING,
amount DECIMAL(10, 2),
is_deleted BOOLEAN,
updated_at TIMESTAMP,
event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(event_ts))
TBLPROPERTIES (
'format-version' = '2',
'write.merge.mode' = 'merge-on-read'
)
""")
import random
from datetime import datetime, timedelta
base_ts = datetime(2024, 1, 1)
statuses = ["created", "paid", "shipped"]
baseline_rows = [
(
order_id,
random.randint(1, 50_000),
random.choice(statuses),
round(random.uniform(5, 500), 2),
False,
base_ts,
base_ts + timedelta(days=order_id % 30, seconds=order_id),
)
for order_id in range(1, 200_001)
]
baseline_df = spark.createDataFrame(
baseline_rows,
schema="""order_id BIGINT, customer_id BIGINT, status STRING,
amount DECIMAL(10,2), is_deleted BOOLEAN,
updated_at TIMESTAMP, event_ts TIMESTAMP""",
)
baseline_df.writeTo("lakehouse.analytics.orders").append()
Это - тот же шаблон Кейса 0 из седьмого урока: создаём v2-таблицу с явно заданным write.merge.mode, и наполняем её 200 000 «исходных» заказов, как если бы это был результат первоначального полного снапшота CDC-коннектора (событие op = 'r' из раздела про событийную модель). Колонка is_deleted добавлена по сравнению со схемой шестого/седьмого уроков специально для этого урока - она реализует паттерн soft delete, разобранный в разделе про WHEN MATCHED THEN DELETE.
Кейс 1: идемпотентный upsert с защитой от устаревших событий и soft delete¶
Сформируем дельту, которая одновременно содержит обновление существующего заказа, вставку нового заказа и логическое удаление третьего - именно такую смесь, какую поставляет реальный CDC-микробатч:
delta_rows = [
# обновление: заказ #42 меняет статус на "shipped"
(42, 17, "shipped", 129.90, False, datetime(2024, 6, 1, 10, 0, 0),
base_ts + timedelta(days=42 % 30, seconds=42)),
# новый заказ, которого не было в baseline
(200_001, 31, "created", 59.00, False, datetime(2024, 6, 1, 10, 0, 0),
datetime(2024, 6, 1)),
# логическое удаление заказа #99 (soft delete, не физический DELETE)
(99, 5, "cancelled", 15.50, True, datetime(2024, 6, 1, 10, 0, 0),
base_ts + timedelta(days=99 % 30, seconds=99)),
]
delta_df = spark.createDataFrame(delta_rows, schema=baseline_df.schema)
delta_df.createOrReplaceTempView("orders_cdc")
spark.sql("""
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc s
ON t.order_id = s.order_id
WHEN MATCHED AND s.updated_at > t.updated_at THEN UPDATE SET
t.status = s.status,
t.amount = s.amount,
t.is_deleted = s.is_deleted,
t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, is_deleted, updated_at, event_ts)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.is_deleted, s.updated_at, s.event_ts)
""")
spark.sql("""
SELECT order_id, status, is_deleted, updated_at
FROM lakehouse.analytics.orders
WHERE order_id IN (42, 99, 200001)
ORDER BY order_id
""").show()
# +--------+--------+----------+-------------------+
# |order_id| status|is_deleted| updated_at|
# +--------+--------+----------+-------------------+
# | 42| shipped| false|2024-06-01 10:00:00|
# | 99|cancelled| true|2024-06-01 10:00:00|
# | 200001| created| false|2024-06-01 10:00:00|
# +--------+--------+----------+-------------------+
Все три строки результата подтверждают теорию по отдельности: заказ 42 обновился, потому что его s.updated_at оказался новее, чем t.updated_at (изначально равный base_ts); заказ 200001 появился через ветку INSERT, потому что не нашёл пары в target; заказ 99 физически остался строкой в таблице (его можно выбрать через SELECT), но получил is_deleted = true - именно так выглядит soft delete на уровне результата: запись не исчезает, а помечается, оставаясь доступной для аналитических запросов, которым нужна полная история, при условии что они сами фильтруют WHERE is_deleted = false.
Проверим идемпотентность - повторно применим тот же самый MERGE INTO с тем же source ещё раз:
snapshot_before = spark.sql("""
SELECT snapshot_id FROM lakehouse.analytics.orders.snapshots
ORDER BY committed_at DESC LIMIT 1
""").collect()[0]["snapshot_id"]
spark.sql("""
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc s
ON t.order_id = s.order_id
WHEN MATCHED AND s.updated_at > t.updated_at THEN UPDATE SET
t.status = s.status,
t.amount = s.amount,
t.is_deleted = s.is_deleted,
t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, is_deleted, updated_at, event_ts)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.is_deleted, s.updated_at, s.event_ts)
""")
second_run_summary = spark.sql("""
SELECT summary['changed-partition-count'] AS changed_partitions
FROM lakehouse.analytics.orders.snapshots
ORDER BY committed_at DESC LIMIT 1
""").collect()[0]["changed_partitions"]
print(second_run_summary)
# 0 - второй прогон не изменил ни одной партиции
changed-partition-count = 0 - формальное, измеримое доказательство идемпотентности: второй прогон того же MERGE INTO с тем же source не находит ни одной строки, удовлетворяющей s.updated_at > t.updated_at (все три заказа теперь имеют t.updated_at, равный s.updated_at из первого прогона, и строгое неравенство не выполняется), поэтому ни одна WHEN MATCHED-ветка не срабатывает, а WHEN NOT MATCHED тоже не находит новых строк, ведь все три заказа уже существуют в target. Это ровно то свойство, которое требуется от production CDC-пайплайна: повторный запуск (например, после переигровки части Kafka-топика для восстановления после сбоя) безопасен и не вносит лишних изменений.
Кейс 2: сравнение физического плана и стоимости MERGE под CoW и MoR¶
Создадим копию таблицы, идентичную по данным, но сконфигурированную под copy-on-write, и применим к обеим таблицам один и тот же микробатч обновлений, чтобы напрямую сравнить результат - не теоретически, а по факту, через системные таблицы и измеренное время:
spark.sql("""
CREATE TABLE lakehouse.analytics.orders_cow
USING iceberg
TBLPROPERTIES ('write.merge.mode' = 'copy-on-write')
AS SELECT * FROM lakehouse.analytics.orders
""")
import random
wide_delta_ids = random.sample(range(1, 200_001), k=5_000)
wide_delta_df = spark.createDataFrame(
[
(oid, 1, "paid", 199.99, False, datetime(2024, 6, 2), base_ts + timedelta(seconds=oid))
for oid in wide_delta_ids
],
schema=baseline_df.schema,
)
wide_delta_df.createOrReplaceTempView("wide_delta")
import time
def run_and_measure(table_name: str) -> dict:
start = time.time()
spark.sql(f"""
MERGE INTO {table_name} t
USING wide_delta s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET
t.status = s.status, t.amount = s.amount, t.updated_at = s.updated_at
""")
elapsed = time.time() - start
summary = spark.sql(f"""
SELECT summary FROM {table_name}.snapshots
ORDER BY committed_at DESC LIMIT 1
""").collect()[0]["summary"]
return {
"elapsed_sec": round(elapsed, 1),
"added_data_files": summary.get("added-data-files"),
"deleted_data_files": summary.get("deleted-data-files"),
"added_delete_files": summary.get("added-delete-files"),
"added_position_deletes": summary.get("added-position-deletes"),
}
print("MoR:", run_and_measure("lakehouse.analytics.orders"))
print("CoW:", run_and_measure("lakehouse.analytics.orders_cow"))
# MoR: {'elapsed_sec': 3.1, 'added_data_files': 1, 'deleted_data_files': 0,
# 'added_delete_files': 1, 'added_position_deletes': 5000}
# CoW: {'elapsed_sec': 11.4, 'added_data_files': 30, 'deleted_data_files': 30,
# 'added_delete_files': 0, 'added_position_deletes': 0}
Результат напрямую подтверждает раздел про физический план выполнения. Под MoR Iceberg добавил ровно один маленький новый data-файл (новые версии 5000 обновлённых строк) и один position delete файл (координаты их старых версий), не тронув ни единого существующего data-файла - и весь MERGE INTO выполнился за 3.1 секунды. Под CoW тот же логический результат потребовал переписать все 30 партиций, потому что случайные 5000 order_id статистически гарантированно затронули хотя бы одну строку в каждой из 30 дневных партиций (тот же эффект «разбросанных изменений», что демонстрировал шестой урок для bucket-партиционирования) - 30 старых файлов помечены deleted-data-files, 30 новых появились как added-data-files, а время выполнения выросло почти в 4 раза, до 11.4 секунды. Открыв Spark UI для обоих запусков, легко увидеть совпадающую длительность стадии Join (одинаковый FULL OUTER JOIN, одинаковая логика сопоставления), но кардинально разную длительность стадии записи (Write) - именно там физически проявляется разница CoW/MoR, а не на уровне самого join'а.
Кейс 3: провокация ошибки cardinality check и её устранение через дедупликацию¶
Сымитируем повторную доставку события из Kafka - источник содержит два разных события по одному и тому же order_id:
duplicate_rows = [
(777, 8, "paid", 75.00, False, datetime(2024, 6, 3, 9, 0, 0), datetime(2024, 6, 3)),
(777, 8, "cancelled", 75.00, True, datetime(2024, 6, 3, 9, 0, 5), datetime(2024, 6, 3)),
]
dirty_df = spark.createDataFrame(duplicate_rows, schema=baseline_df.schema)
dirty_df.createOrReplaceTempView("orders_cdc_dirty")
try:
spark.sql("""
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc_dirty s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET
t.status = s.status, t.is_deleted = s.is_deleted, t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, is_deleted, updated_at, event_ts)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.is_deleted, s.updated_at, s.event_ts)
""")
except Exception as e:
print(type(e).__name__, "-", str(e)[:160])
# SparkException - ... a single row from the target table was matched
# with multiple rows of the source table ...
Ошибка возникает ровно там, где предсказывает теория: order_id = 777 встречается в orders_cdc_dirty дважды, оба раза проходят ON-условие и сопоставляются с одной и той же строкой target, и cardinality check прерывает выполнение прежде, чем хоть одна строка успела бы измениться. Проверим это явно:
unaffected = spark.sql("""
SELECT order_id, status FROM lakehouse.analytics.orders WHERE order_id = 777
""").collect()[0]
print(unaffected)
# Row(order_id=777, status='created') - исходный статус, MERGE откатился целиком
Применим дедупликацию по ROW_NUMBER() из теоретического раздела и повторим попытку:
spark.sql("""
CREATE OR REPLACE TEMP VIEW orders_cdc_deduped AS
SELECT order_id, customer_id, status, amount, is_deleted, updated_at, event_ts
FROM (
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY order_id ORDER BY updated_at DESC
) AS rn
FROM orders_cdc_dirty
)
WHERE rn = 1
""")
spark.sql("""
MERGE INTO lakehouse.analytics.orders t
USING orders_cdc_deduped s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET
t.status = s.status, t.is_deleted = s.is_deleted, t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, status, amount, is_deleted, updated_at, event_ts)
VALUES (s.order_id, s.customer_id, s.status, s.amount, s.is_deleted, s.updated_at, s.event_ts)
""")
spark.sql("""
SELECT order_id, status, is_deleted FROM lakehouse.analytics.orders WHERE order_id = 777
""").show()
# +--------+---------+----------+
# |order_id| status|is_deleted|
# +--------+---------+----------+
# | 777|cancelled| true|
# +--------+---------+----------+
После дедупликации MERGE INTO применяется без единой ошибки, и в target оказывается ровно та версия события, у которой updated_at максимален (cancelled, 09:00:05) - именно «победитель» из двух конкурирующих дублей по правилу ORDER BY updated_at DESC. Этот кейс демонстрирует полный цикл production-инцидента в миниатюре: воспроизведение ошибки, диагностика её причины через текст исключения и состояние target, и устранение через дополнительный CTE-слой дедупликации перед USING-секцией, без единого изменения самой структуры WHEN-веток.
Производственный кейс: тысячи мелких файлов от высокочастотного MERGE INTO¶
Команда платформы данных одного из продуктов e-commerce построила потоковый пайплайн репликации заказов: Debezium читает WAL продуктовой PostgreSQL, Kafka доставляет события, а Spark Structured Streaming job с trigger(processingTime="1 minute") каждую минуту материализует микробатч в стейджинг-таблицу и применяет к Iceberg-таблице lakehouse.analytics.orders MERGE INTO, построенный по той же схеме, что разобрана в Кейсе 1 этого урока: дедупликация по ROW_NUMBER(), guard по updated_at, soft delete через is_deleted. Таблица была сконфигурирована с write.merge.mode = 'merge-on-read' - решение, полностью обоснованное материалом этого урока: высокочастотный поток обновлений, разбросанных по случайным заказам, через CoW обходился бы кратно дороже, как показал Кейс 2.
Первые недели всё работало штатно: каждый MERGE INTO коммитился за 2-4 секунды, что комфортно укладывалось в минутный интервал триггера. Но команда не настроила write.distribution-mode явно, оставив значение по умолчанию none - и не считала это значимым решением, так как раздел про этот параметр на момент запуска пайплайна ещё не был изучен никем в команде.
Диагностика подтвердила механизм, разобранный в разделе про write.distribution-mode: при none каждая executor-задача каждого минутного микробатча писала данные в любые партиции, которые встретились среди строк, которые ей досталось обработать - а поскольку CDC-события по своей природе разбросаны по случайным event_ts, почти каждый микробатч создавал собственный мелкий файл в каждой партиции, которую он касался, вместо того чтобы консолидировать запись в одну партицию через одну задачу. За четыре недели непрерывной работы 1440 микробатчей в сутки (один в минуту) суммарно создали 48 000 data-файлов - притом что логический объём данных был эквивалентен таблице, которая при батчевой загрузке раз в сутки потребовала бы не более 60-90 файлов. Рост числа manifest-файлов, обслуживающих эти data-файлы, напрямую увеличил объём метаданных, которые Scan Planning (третий урок модуля) должен прочитать и обработать для любого запроса к таблице, независимо от того, насколько узкую партицию этот запрос затрагивает - именно поэтому пострадали даже точечные BI-запросы, логически не нуждавшиеся в большей части истории таблицы.
| Метрика | Неделя 1 | Неделя 4 (пик) | После исправления |
|---|---|---|---|
| Data-файлов в таблице | ~1 500 | 48 000 | ~4 200 |
| Manifest-файлов | ~60 | 1 200+ | ~85 |
| Средний размер data-файла | 4.1 МБ | 340 КБ | 38 МБ |
| Latency типового BI-запроса | ~1.5 сек | 9-12 сек | ~2.1 сек |
Меры по итогам инцидента:
-
Немедленная мера: на таблицу включена настройка
write.distribution-mode = 'hash', гарантирующая, что каждая партиция получает данные ровно от одной (или предсказуемо малой группы) executor-задачи на каждый микробатч, вместо разбросанной по множеству задач записи. -
Структурное исправление: запущена компакция через
rewrite_data_files(детальный синтаксис - тема десятого урока этого модуля) для немедленного сведения существующих 48 000 файлов к компактному набору, приближенному к целевому размеру файла. -
Архитектурное уточнение: команда зафиксировала правило для всех новых таблиц с высокочастотным CDC-потоком -
write.distribution-modeобязателен к явной настройке при создании таблицы (а не оставляется на дефолт «как получится»), а количество data- и manifest-файлов добавлено в стандартный дашборд мониторинга наравне с latency самогоMERGE INTO, чтобы рост числа файлов был виден заранее, а не постфактум через деградацию чтения.
Важная связь с материалом следующего раздела про типичные заблуждения: этот инцидент не является аргументом против merge-on-read как такового - сам выбор режима был обоснован профилем нагрузки. Корневая причина - отдельная, независимая настройка (write.distribution-mode), которую легко упустить, считая, что единственный относящийся к производительности параметр MERGE INTO - это write.merge.mode. На практике оба рычага конфигурации работают совместно и должны настраиваться согласованно.
Типичные заблуждения¶
«MERGE INTO - это просто синтаксический сахар над отдельными UPDATE/DELETE/INSERT» - неверно и упускает главную причину существования оператора. Раздел про атомарность микробатча показал, что именно неделимость коммита всей смеси изменений за одну операцию решает задачу, которую три отдельные команды решить не способны без промежуточных несогласованных состояний таблицы.
«MERGE INTO всегда генерирует и position, и equality delete-файлы одновременно» - неточная трактовка, требующая уточнения относительно материала седьмого урока. Когда MERGE INTO выполняется как обычная SQL-команда внутри Spark (ровно тот сценарий, что разобран в этом уроке), движок всегда выполняет собственный join/scan и поэтому всегда точно знает позиции совпадающих строк target - то есть, как и для plain UPDATE/DELETE, физическим результатом всегда оказываются position deletes, а не equality. Формулировка «MERGE INTO генерирует оба типа delete-файлов» относится не к самому SQL-оператору, а к более широкому архитектурному понятию «upsert-паттерн», которое на практике иногда реализуется не через SQL MERGE INTO, а через внешний потоковый writer (например, кастомный Flink-CDC sink, разобранный в седьмом уроке), использующий низкоуровневое Java API Iceberg в режиме blind write - именно такой writer, а не сам оператор MERGE INTO, порождает equality deletes.
«Чем больше WHEN-веток, тем дороже запрос» - неверно как общее правило. Раздел про физический план показал, что стоимость определяется не количеством веток, а тем, какие типы веток присутствуют (это определяет тип join) и какой объём данных реально требует физической перезаписи (CoW) или добавления (MoR) - можно добавить произвольное число дополнительных предикатов внутри уже существующих WHEN MATCHED/WHEN NOT MATCHED без изменения типа join вообще, и стоимость от этого практически не меняется.
«Cardinality check - это просто перестраховка, которую можно отключить, если уверены в данных» - рискованная позиция. write.merge.cardinality-check.enabled = false отключает защитный механизм, но не устраняет саму неоднозначность: при наличии дублей в source поведение становится недетерминированным (порядок применения двух конкурирующих обновлений к одной строке target не гарантирован), просто без явной ошибки, сигнализирующей об этом. Отключение проверки превращает явную, диагностируемую ошибку в скрытую, проявляющуюся как необъяснимая нестабильность данных.
«ROW_NUMBER()-дедупликация решает все проблемы консистентности CDC» - неполно. Раздел про late-arriving data показал, что внутрибатчевая дедупликация устраняет только конфликты внутри одного микробатча; защита от опоздавших событий между разными микробатчами требует отдельного guard-условия по updated_at/версии, и обе защиты необходимы одновременно, а не взаимозаменяемы.
«write.distribution-mode = hash всегда улучшает производительность MERGE INTO» - неверно как общее правило, прямо обсуждённое в разделе про эту настройку: дополнительный shuffle перед записью имеет собственную цену, и для очень маленьких микробатчей эта цена может превышать выгоду от консолидации файлов. Производственный кейс этого урока показывает реальный сценарий, где hash помог, но решение должно приниматься на основе измерения конкретного профиля нагрузки, а не как универсальное правило «всегда включать».
«Soft delete (is_deleted) - это просто альтернативная реализация DELETE, выбор между ними не имеет значения» - неточно. Раздел про обработку удалений показал, что это архитектурное решение влияет на видимость исторических состояний строки для аналитики и аудита, а не только на физический механизм записи - выбор должен делаться осознанно на уровне бизнес-требований к таблице, а не как деталь реализации MERGE INTO.
Производственный чек-лист: что нельзя забыть до включения CDC-пайплайна в production¶
Сжатый итог решений, разобранных в этом уроке, в форме контрольного списка, которым стоит сверяться при проектировании новой CDC-таблицы:
-
Бизнес-ключ в
ON-условии действительно уникален в целевой таблице - проверено явно (например, черезSELECT order_id, count(*) FROM target GROUP BY order_id HAVING count(*) > 1), а не предполагается по дизайну схемы источника. -
Source перед
USINGдедуплицирован черезROW_NUMBER() OVER (PARTITION BY <key> ORDER BY <freshness>)- на постоянной основе, как часть пайплайна, а не как разовое исправление после первого столкновения с cardinality check. -
Guard-условие по
updated_at(или аналогичной версии/LSN) присутствует во всехWHEN MATCHED-ветках, изменяющих данные - защита от опоздавших и переигранных событий встроена в сам запрос, а не полагается на гарантии порядка доставки источника. -
Партиционирование таблицы выбрано осознанно относительно ключа join
MERGE INTO, а не только относительно типичных аналитических запросов - раздел про перекос данных показал, что несогласованность этих двух измерений напрямую создаёт data skew, который пруннинг по партициям не способен компенсировать. -
write.merge.modeуказан явно, а не оставлен на дефолтное значение - выбор между CoW и MoR должен быть результатом анализа профиля нагрузки (частотаMERGE INTO, доля изменённых строк, требования к latency чтения), а не случайностью версии Iceberg по умолчанию. -
write.distribution-modeнастроен согласованно с частотой и шириной микробатчей - производственный кейс этого урока показал, что отсутствие этой настройки превращает высокочастотный CDC-поток в генератор мелких файлов независимо от того, насколько правильно выбранwrite.merge.mode. -
Maintenance-процедура компакции запланирована заранее, а не добавляется реактивно после первого инцидента деградации чтения - это прямое продолжение урока про Read Penalty (седьмой урок) и центральная тема следующего урока модуля.
Последний пункт чек-листа заслуживает отдельного, более развёрнутого комментария, потому что именно он формирует мостик к следующей теме модуля. MERGE INTO в режиме merge-on-read, запускаемый часто (что типично для production CDC-пайплайна с интервалом триггера от секунд до нескольких минут), системно накапливает технический долг в виде растущего числа мелких data-файлов и delete-файлов - даже при идеально настроенном write.distribution-mode, потому что сама природа потоковой нагрузки (много мелких независимых коммитов вместо одной крупной batch-загрузки) механически означает много мелких файлов на каждый коммит. Это - не ошибка конфигурации и не повод отказываться от merge-on-read, а неизбежное следствие самой архитектуры построчного потокового upsert'а, требующее системного, регулярного противодействия через компакцию - именно тот процесс, который подробно разбирает следующий урок модуля.
Мостик к следующим урокам¶
Этот урок завершил трилогию построчных операций модуля (Copy-on-Write → Merge-on-Read → MERGE INTO), показав, как обе физические стратегии записи объединяются в единый production-паттерн upsert'а для CDC-репликации. Следующие уроки модуля строят непосредственно на материале этого урока.
-
Урок 9 («Time Travel») покажет, как обращаться к историческим снапшотам таблицы, прошедшей через десятки и сотни коммитов
MERGE INTO- в том числе как корректно интерпретировать состояние строки на момент времени до того, как конкретное CDC-событие было применено, опираясь на тот же sequence number, который определяет порядок применения delete-файлов из седьмого урока. -
Урок 10 («Table Maintenance») - прямое и самое непосредственное продолжение последнего раздела этого урока: производственный кейс и чек-лист выше многократно анонсировали
rewrite_data_filesиexpire_snapshotsкак обязательную, а не опциональную часть архитектуры любой таблицы с активным CDC-потоком черезMERGE INTO. Этот урок разберёт точный синтаксисCALL-процедур, стратегии планирования компакции (по расписанию, по порогу числа файлов, по доле мелких файлов) и то, какexpire_snapshotsбезопасно освобождает физическое место, занятое старыми версиями файлов, переписанных предыдущимиMERGE INTO. -
Уроки 11-12 («Delta Lake») покажут, как тот же паттерн upsert для CDC реализован в альтернативном table format - оператор
MERGE INTOсинтаксически почти идентичен, но физическая модель отличается: вместо дерева манифестов Delta Lake использует transaction log (Delta Log), а вместо отдельных position/equality delete files - механизм Deletion Vectors, решающий концептуально ту же задачу «отметить строку как удалённую без переписывания файла», но иным техническим устройством. -
Урок 13 («Выбор формата») включит производительность и эргономику
MERGE INTOкак один из критериев сравнительной матрицы Iceberg vs Delta Lake vs Apache Hudi - в частности, то, как именно каждый формат балансирует стоимость записи и стоимость чтения для одного и того же CDC-профиля нагрузки, разобранного в этом уроке.
Домашнее задание¶
-
Воспроизведите Кейс 1 практического блока на собственном self-hosted стенде, но добейтесь срабатывания cardinality check не через явные дубли в Python-списке (как в Кейсе 3), а через реалистичный сценарий: смоделируйте at-least-once повторную доставку, добавив в
orders_cdcкопию одной случайной строки с тем жеorder_id, но другимupdated_at, через DataFrame-операциюunion. Убедитесь, что ошибка воспроизводится, и письменно опишите, по какому именно полю текста исключения можно было бы построить автоматический алерт. -
Постройте таблицу (текстом или через matplotlib) зависимости времени выполнения
MERGE INTOот доли строк целевой таблицы, затронутых обновлением (0.1%, 1%, 5%, 25%, 50%), отдельно дляwrite.merge.mode = 'copy-on-write'и'merge-on-read', используя подход Кейса 2. Определите экспериментально приблизительную долю изменённых строк, при которой два режима показывают сопоставимое время выполнения, и объясните результат через формулу Write Amplification из шестого урока. -
Реализуйте паттерн
WHEN NOT MATCHED BY SOURCE(раздел про стандарт ANSI) как явную дополнительную ветку через подзапросEXCEPT, если ваша версия Iceberg-Spark runtime не поддерживает синтаксис нативно: найдите заказы, присутствующие в target, но отсутствующие в полном дневном снапшоте источника, и пометьте ихis_deleted = trueотдельной командойUPDATE. Письменно сравните этот подход с нативнымMERGE INTOпо читаемости и по числу операций, требующих отдельного коммита. -
Спроектируйте (без необходимости полной реализации) Structured Streaming
foreachBatch-функцию, которая применяет паттерн дедупликации и guard-условие этого урока к каждому микробатчу автоматически, принимая имя целевой таблицы и список колонок бизнес-ключа как параметры. Обсудите письменно, какие дополнительные проверки (помимо cardinality check самого Iceberg) стоило бы добавить перед вызовомMERGE INTOвнутри такой функции. -
Используя дерево решений join-стратегии из раздела про физический план, определите письменно, какой логический join будет построен для запроса, содержащего только
WHEN MATCHED AND s.op = 'd' THEN DELETE(без единой веткиUPDATEилиINSERT), и объясните, почему результат отличается от случаяWHEN MATCHED THEN UPDATE, хотя оба - «только WHEN MATCHED» с точки зрения дерева решений. -
Воспроизведите производственный кейс этого урока на собственном стенде: создайте таблицу без явного
write.distribution-mode, запустите 50 последовательныхMERGE INTOсо случайно разбросанными по партициям микробатчами по 100 строк, замерьте число файлов через.filesдо и после, затем повторите эксперимент сwrite.distribution-mode = 'hash'с нуля. Письменно сравните итоговое число файлов и объясните разницу через механизм, разобранный в разделе про эту настройку. -
Настройте Bloom-фильтр для
customer_idна тестовой таблице из 1-2 млн строк (write.parquet.bloom-filter-enabled.column.customer_id = 'true'), выполнитеMERGE INTOсON t.customer_id = s.customer_idдля микробатча из нескольких сотен строк, и сравните измеренное время выполнения с тем жеMERGE INTOна копии таблицы без включённого Bloom-фильтра. Объясните полученный (или отсутствующий) выигрыш через рассуждение об избирательности фильтра и распределенииcustomer_idв данных. -
Напишите unit-тест (PySpark +
pytest, локальный Iceberg-каталог), проверяющий инвариант идемпотентности из Кейса 1: повторное применение одного и того жеMERGE INTOс одним и тем же source два раза подряд должно датьchanged-partition-count = 0для второго прогона. Тест должен явно покрывать все три типа событий (INSERT,UPDATE, softDELETE) в одном source. -
Спроектируйте письменно (1-2 абзаца) конкретную метрику и пороговое значение для алертинга, которые позволили бы обнаружить деградацию из производственного кейса этого урока (рост числа файлов) на неделе 1, а не на неделе 4 - аналогично тому, как был спроектирован алерт в производственном кейсе седьмого урока, но используя другой источник данных (
.files, а не.delete_files). -
Сравните письменно (без необходимости полной реализации) то, как изменился бы дизайн пайплайна этого урока, если бы целевая таблица не допускала никакой задержки между событием в источнике и его видимостью в Iceberg-таблице (триггер раз в несколько секунд вместо раз в минуту). Какие из решений урока (дедупликация, guard-условие,
write.distribution-mode) продолжают работать без изменений, а какие требуют пересмотра, и почему - опираясь на разницу между «логической корректностью» и «физической стоимостью» каждого механизма.
Полная картина: жизненный цикл одного CDC-микробатча от события до компакции¶
Завершая урок, соберём весь материал в единую диаграмму - от события в исходной СУБД до момента, когда накопленные малые файлы наконец компактизируются, замыкая цикл на тему следующего урока модуля.
Эта диаграмма связывает воедино весь материал урока в порядке, в котором реальный CDC-микробатч проходит через систему: от транзакции в источнике (раздел про событийную модель) через дедупликацию и guard-условие (раздел про каверзные сценарии) к выбору join-стратегии и физическому режиму записи (раздел про физический план), заканчивая атомарным коммитом и - при недостаточном контроле - постепенным накоплением технического долга в виде мелких файлов, разрешаемым только регулярной компакцией (производственный кейс и мостик к десятому уроку). Каждый блок этой схемы - не абстракция, а конкретный, измеренный в практическом блоке этого урока механизм.
Итоги¶
MERGE INTO - единственный из трёх инструментов записи Iceberg, принимающий именно то, что поставляет CDC-источник: дельту, а не полный снимок. В отличие от INSERT OVERWRITE (требует полного состояния) и APPEND (накапливает дубли по ключу), MERGE INTO за одну атомарную операцию сопоставляет источник с целевой таблицей и применяет вставку, обновление и удаление одновременно.
Синтаксис MERGE INTO следует стандарту ANSI SQL и состоит из трёх типов веток, оцениваемых по порядку. WHEN MATCHED THEN UPDATE/DELETE и WHEN NOT MATCHED THEN INSERT ведут себя как ветки CASE WHEN; если строка MATCHED, но не satisfies ни одного условия - она остаётся нетронутой без ошибки, что и делает возможным безопасный guard против устаревших событий.
Логический тип join, который строит планировщик, определяется набором присутствующих WHEN-веток, а не настройками конфигурации. Только WHEN MATCHED - RIGHT OUTER JOIN; WHEN MATCHED + WHEN NOT MATCHED - FULL OUTER JOIN; только WHEN NOT MATCHED - LEFT ANTI JOIN. Этот логический план одинаков для CoW и MoR; различается только физическая запись результата.
Под CoW MERGE INTO обязан переписать весь затронутый файл целиком, включая несовпавшие строки; под MoR - только записать координаты изменённых строк и маленький новый файл, не трогая остальное содержимое старого файла. Это - прямое продолжение моделей Read-Modify-Write и LSM-подобного журнала изменений из шестого и седьмого уроков, применённое теперь к комбинации UPDATE/DELETE/INSERT за одну операцию.
Три table property независимо настраивают производительность: write.merge.mode выбирает физическую стратегию, write.distribution-mode управляет числом и размером итоговых файлов через shuffle перед записью, Bloom-фильтры ускоряют точечный поиск совпадений внутри уже отобранных файлов. Все три рычага дополняют друг друга, а не заменяют - производственный кейс урока показал, что упущение write.distribution-mode создаёт проблему мелких файлов независимо от того, насколько верно выбран write.merge.mode.
Дубли в источнике и опоздавшие события - не исключение, а норма для CDC-потоков, и Iceberg явно защищается от первой проблемы через cardinality check. Решение - дедупликация по ROW_NUMBER() внутри батча и guard-условие по updated_at/версии между батчами; обе защиты обязательны одновременно, поскольку решают разные, не перекрывающиеся классы ошибок.
Идемпотентность MERGE INTO - проверяемое, а не предполагаемое свойство. Повторный прогон одного и того же запроса с тем же source при правильно спроектированном guard-условии должен давать changed-partition-count = 0 - именно так была формально подтверждена идемпотентность в практическом блоке урока.
Частый MERGE INTO под MoR системно накапливает технический долг в виде растущего числа мелких файлов - это неизбежное следствие архитектуры потокового upsert'а, а не признак ошибки конфигурации, и требует регулярной, заранее спланированной компакции. Этот тезис - прямой мостик к десятому уроку модуля, где разбирается синтаксис и стратегии планирования rewrite_data_files/expire_snapshots.
Краткий глоссарий терминов урока¶
-
CDC (Change Data Capture) - класс технологий, фиксирующих и транслирующих построчные изменения в исходной СУБД (через журнал транзакций - WAL/binlog) в виде структурированного потока событий, без дополнительных
SELECT-запросов к самому источнику. -
Upsert - паттерн «обновить существующую строку или вставить новую, если она не существует», реализуемый в Iceberg через
MERGE INTO; часто дополняется обработкой удалений (Upsert + Delete). -
MERGE INTO- SQL-оператор, сопоставляющий целевую (target) и исходную (source) таблицы по условиюONи применяющий вставку, обновление или удаление в зависимости от того, какаяWHEN-ветка сработала первой для каждой строки. -
WHEN MATCHED/WHEN NOT MATCHED- веткиMERGE INTO, оцениваемые по порядку какCASE WHEN; первая - для строк, нашедших пару поON-условию, вторая - для строк source без пары в target. -
Cardinality check - встроенная проверка Iceberg (
write.merge.cardinality-check.enabled), прерывающая весьMERGE INTO, если одна строка target оказалась сопоставлена более чем с одной строкой source в рамкахWHEN MATCHED. -
write.merge.mode- table property, выбирающая физическую стратегию выполненияMERGE INTO:copy-on-write(переписывание затронутых файлов целиком) илиmerge-on-read(добавление delete- и маленьких data-файлов). -
write.distribution-mode- table property (none/hash/range), управляющая тем, выполняется ли shuffle данных по ключу партиционирования перед физической записью, напрямую влияющая на число и размер итоговых файлов. -
Bloom-фильтр (Parquet) - вероятностная структура данных уровня row group, отвечающая «может ли значение присутствовать» без false negative; ускоряет точечный поиск совпадений по ключу join внутри уже отобранных Scan Planning'ом файлов.
-
Дедупликация по
ROW_NUMBER()- идиоматичный паттерн устранения дублей источника по бизнес-ключу передMERGE INTO, оставляющий ровно одну, самую свежую версию каждой строки черезPARTITION BY <key> ORDER BY <freshness> ... WHERE rn = 1. -
Late-arriving data - событие источника, доставленное после того, как более новое состояние той же сущности уже было применено к target; защищается guard-условием вида
s.updated_at > t.updated_atвнутриWHEN MATCHED. -
Идемпотентность MERGE - свойство, при котором повторное применение одного и того же
MERGE INTOс тем же source не вносит дополнительных изменений в target; формально проверяется черезchanged-partition-count = 0второго прогона. -
Data Skew (перекос данных) - неравномерное распределение значений ключа join, приводящее к тому, что часть executor-задач
MERGE INTOполучает кратно больше строк, чем остальные, и становится «стрэглером», определяющим общее время выполнения операции. -
Soft delete - паттерн логического удаления, при котором строка не исчезает физически из таблицы, а помечается флагом (например,
is_deleted = true) через веткуWHEN MATCHED THEN UPDATE, сохраняя видимость исторического состояния для аудита и аналитики. -
Staging-таблица CDC - промежуточная таблица (Iceberg или временная), в которую микробатч потоковых событий материализуется перед тем, как стать source для
MERGE INTO, упрощая инспекцию, дедупликацию и переигровку дельты. -
op(поле CDC-события) - стандартизированное поле события Debezium, кодирующее тип изменения:c(create/insert),u(update),d(delete),r(read, первоначальный snapshot); используется вMERGE INTOдля выбора правильнойWHEN-ветки.