Iceberg и Paimon — строительные блоки Streaming Lakehouse


Привет, Хабр! Я Алексей Новаков, ведущий инженер данных в Рунити. Сегодня разберемся, зачем Lakehouse понадобился streaming, почему одного Iceberg для этого не всегда хватает и как Iceberg и Paimon могут жить рядом в одном конвейере.
Если сильно упростить, классический Lakehouse решает понятную задачу: берем дешевое объектное хранилище и добавляем поверх него то, чего не хватало обычному Data Lake — транзакции, schema evolution, snapshots, SQL, update/delete отдельных строк и одновременное чтение и запись.
Но дальше возникает следующий вопрос: что делать, если аналитика должна видеть не вчерашние или часовые данные, а состояние, которое обновляется постоянно?
Именно здесь Lakehouse начинает превращаться в Streaming Lakehouse — или Streamhouse.
Навигация по тексту:
Таблица — это уже не просто набор Parquet-файлов
Сами данные в Lakehouse по-прежнему могут лежать в Parquet или ORC. Но одного колоночного файла мало: кто-то должен хранить схему, версии таблицы, список актуальных файлов, snapshots и информацию для транзакций.
Эту работу берут на себя открытые табличные форматы — Iceberg, Paimon, Hudi, Delta Lake. В архитектуре Lakehouse они становятся прослойкой между движком обработки и физическим хранилищем вроде S3.
Iceberg таблица физически выглядит примерно так:
s3://lakehouse/iceberg_catalog/table_1/
├── metadata/
│ ├── *.json
│ └── *.avro
└── data/
├── partition-0/
│ └── *.parquet
└── partition-N/
└── *.parquet
Метаданные описывают snapshots, manifests и актуальный набор data-файлов. Сами данные остаются в Parquet. В этом примере мы также используем партиционирование по одному из столбцов таблицы.
Paimon устроен немного иначе, но суть метаинформации в нем такая же:
s3://lakehouse/paimon_catalog/table_2/
├── snapshot/
├── schema/
├── manifest/
├── bucket-0/
│ └── data-*.orc
└── bucket-N/
└── data-*.orc
Внутри bucket'ов данные организованы как SST-файлы в LSM-структуре. И вот это различие становится особенно заметным, когда вместо периодической batch-обработки начинаются постоянные изменения данных.
Где у Iceberg появляются неудобства
Iceberg хорош как универсальный табличный формат: у него много интеграций, с ним работают Flink, Spark, Trino и другие инструменты. С append-only потоком всё тоже достаточно понятно. События приходят, Flink пишет новые данные в таблицу, downstream-job читает новые snapshots.
Проблема начинается с upsert-таблиц.
Представим обычный pipeline:
загрузка → трансформация → агрегация → витрина
На первом этапе события можно просто добавлять. Но после агрегации или CDC появляется вопрос о последнем состояние конкретной записи: одну и ту же запись нужно обновлять по primary key.
Iceberg поддерживает upsert, но дальше возникает ограничение: streaming consumer не получает из такой таблицы полноценную историю (changelog) изменений. То есть использовать одну upsert-таблицу как источник следующего непрерывного upsert/append-конвейера будет невозможно, вы будите получать только новые записи, но не обновления. Один из вариантов обхода — снова выносить changelog в Kafka, то есть создать промежуточную таблицу поверх Kafka топика.
И тут появляется Paimon.
Он изначально гораздо сильнее заточен под связку с Flink и потоковые изменения:
append → append;
append → upsert;
upsert → upsert;
upsert → append;
streaming read для changelog;
temporal lookup join;
Flink SQL с полной поддержкой watermarks конфигурации.
Поэтому на практике вопрос обычно не звучит как «Iceberg или Paimon?». Гораздо интереснее подобрать формат под конкретный тип таблицы.
Собираем один пайплайн из двух форматов
Для демо я взял сценарий аналитики маркетинговых кампаний.
Есть два потока данных:
Первый — web events. Пользователь открыл страницу, кликнул, сделал conversion. Это естественный append-only поток.
Второй — данные рекламных кампаний. Campaign name, status, customer segment, UTM и остальные атрибуты. Эти данные иногда меняются, но в исходной таблице немного строк, поэтому копируем эти все рекламные кампании каждый раз полностью в новую партицию.
Получилась такая схема:

Физически все это может жить в одном S3-бакете:
s3://lakehouse/
├── iceberg/
│ ├── bronze/web_events/
│ ├── silver/enriched_events/
│ └── gold/campaign_metrics/
├── paimon/
│ └── marketing.db/campaigns/
├── flink-checkpoints/
└── flink-savepoints/
То есть Iceberg и Paimon не требуют двух отдельных хранилищ типа S3 или HDFS. Они спокойно живут рядом на одном объектном хранилище, каждый решая свою часть задачи.
События пишем в Iceberg
Web events просто складываются в bronze:
INSERT INTO iceberg_rest.bronze.web_events
SELECT
event_id,
event_type,
user_id,
session_id,
page_url,
utm_source,
utm_medium,
utm_campaign,
event_timestamp
FROM web_events_datagen;
В демо таким образом накопилось около 220 тысяч событий.
Справочник кампаний — в Paimon
Кампании лежат в Paimon как изменяемое состояние — в примере их около 50.
При обогащении web events можно сделать temporal lookup:
LEFT JOIN paimon_catalog.marketing.campaigns
/*+ OPTIONS(
'scan.partitions'='max_pt()',
'lookup.dynamic-partition.refresh-interval'='1 h'
) */
FOR SYSTEM_TIME AS OF b.proc_time AS m
ON b.utm_source = m.utm_source
AND b.utm_campaign = m.utm_campaign
Поток событий при этом не нужно останавливать или периодически перезагружать целиком. Flink берет актуальное состояние справочника и обогащает новые события. После этого данные уходят в silver-слой.
В gold считаем метрики каждую минуту
На последнем шаге Flink агрегирует события по окнам:
FROM CUMULATE(
TABLE silver_enriched,
DESCRIPTOR(event_timestamp),
INTERVAL '1' MINUTE,
INTERVAL '1' HOUR
)
В результате gold.campaign_metrics получает новую статистику каждую минуту: количество событий, пользователей, сессий, просмотров, кликов и конверсий.
Почему это всё-таки не замена Kafka для любого real-time
В нашем демо общая задержка получилась около 4–5 минут.
Примерно три минуты дают последовательные checkpoints в web_events, enriched_events и campaign_metrics, еще около минуты — окно агрегации.
То есть по latency такая схема проигрывает классическому streaming через Kafka. Но это сознательный компромисс: данные сразу оказываются в долговечном табличном хранилище, доступны для SQL и исторического анализа и не требуют отдельного пути из стрима в Lakehouse.
Поэтому я бы смотрел на выбор примерно так:
нужна минимальная задержка и событие имеет высокую бизнес-ценность → классический Flink + Kafka;
данные можно обновлять с задержкой в минуты, но хочется единого долговечного слоя → Streaming Lakehouse;
задержка еще менее критична → обычный batch Lakehouse.
Streamhouse здесь не «новая архитектура, которая отменяет Kafka», а еще одна точка на шкале между latency, стоимостью и сложностью эксплуатации.
А потом приходят снэпшоты и тысячи мелких файлов
На демо всё красиво. В проде появляются две вещи, которые быстро заставляют пересмотреть подход: Consumer может отстать от snapshots.
Streaming reader последовательно идет по snapshots таблицы. Если maintenance слишком быстро удаляет историю, consumer может еще читать старый snapshot, а его уже пора чистить. В Paimon для этого есть consumer свойство таблицы: все клиенты таблицы видят что есть конкретные consumers и их прогресс относительно номера snapshot, и поэтому основной потребитель данных может не бояться что не сможет обработать все изменения. При этом процесс очистки snapshot будет ожидать пока потребитель не перейдет на следующие snapshots. Прогресс потребителей можно сбрасывать или переустанавливать на нужный snapshot_id. То есть lifecycle snapshots в streaming-сценарии уже нельзя настраивать независимо от читателей. Этот момент напоминаем Kafka API с её topic consumer group offsets.
Streaming плодит маленькие файлы
Частые комиты, постоянные и обновления легко создают тысячи и десятки тысяч файлов в день. А дальше растут время листинга, стоимость сканирования, количество операций с object storage и накладные расходы при чтении и записи. Поэтому compaction — не необязательная оптимизация «на потом», а часть самого streaming-процесса.
У Iceberg для этого есть RewriteDataFiles, ExpireSnapshots, DeleteOrphanFiles и post-commit compaction. У Paimon — compact, compact_database и асинхронное слияние файлов.
Что сразу заложить для продакшена
Во-первых, включить file compaction прямо в streaming jobs.
Во-вторых, отдельно запустить maintenance для Iceberg: rewrite data files, очистку orphan files и устаревших snapshots.
В-третьих, заранее автоматизировать восстановление Flink state. Нужно понимать, что делать с поврежденным checkpoint: откуда перечитывать данные, какой Iceberg snapshot использовать и где потребуется idempotent upsert.
В-четвертых, для Kubernetes-развертывания использовать Flink Kubernetes Operator, а не собирать жизненный цикл jobs вручную.
Вывод
Главное, что мне нравится в этой архитектуре — не приходится выбирать один табличный формат на всё озеро данных. Iceberg хорошо подходит для большого количества обычных Lakehouse-таблиц, интегрируется с разными движками и удобно работает с append-oriented слоями. Paimon можно поставить туда, где начинается настоящий streaming state: CDC, upsert, changelog и lookup по постоянно меняющимся данным. А Flink становится движком, который соединяет оба мира в один pipeline.
В демо мы получили полностью потоковый bronze → silver → gold конвейер на S3 с задержкой порядка 4–5 минут. Не быстрее Kafka — и в этом нет цели. Зато одни и те же таблицы одновременно становятся и состоянием streaming-пайплайна, и долговечным слоем данных, доступным для дальнейшей аналитики. Для части задач этого компромисса вполне достаточно.
KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.