Iceberg Sink и Confluent Tableflow
В 2025 году streaming и lakehouse окончательно сошлись в одной архитектуре. Kafka хранит события несколько часов или дней (retention), Iceberg/Delta Lake — годами. Между ними нужен мост: компонент, который материализует поток сообщений в колоночные файлы object-storage и регистрирует их в каталоге как таблицу. Этот мост строится двумя способами: open-source Apache Iceberg Sink Connector для Kafka Connect и управляемый сервис Confluent Tableflow.
Этот урок объясняет оба подхода, их commit-семантику, schema evolution и production-tradeoffs.
Контекст: Streaming-Lakehouse convergence
Классическая архитектура разделяла операционный поток (Kafka) и аналитическое хранилище (Snowflake, BigQuery, Spark). Между ними стояли batch ETL-задачи: каждый час Spark-job читал Kafka, дедуплицировал, писал в Parquet. Лаг — часы. Дублирование данных — двойное хранилище.
В 2024–2025 годах переход на open table formats (Apache Iceberg, Delta Lake, Apache Hudi) изменил картину. Iceberg-таблица — это не просто Parquet-файлы, а транзакционная абстракция со snapshot-isolation и атомарными commit. Это позволило построить streaming ingestion: коннектор пишет файлы прямо в object-storage и фиксирует snapshot, как только пакет готов.
Apache Iceberg по своей природе таблица: каждая строка имеет первичный ключ, schema-evolution встроена в спецификацию, snapshot изолированы. Kafka — упорядоченный лог. Соединение этих абстракций (Kafka topic как Iceberg table) — главный архитектурный паттерн 2025 года, который Jack Vanlightly называет “stream/table duality”.
Old: Kafka + Batch ETL
Классический подход: Kafka хранит события 7 дней, batch ETL раз в час читает их в Spark, дедуплицирует, пишет в Parquet/Hive. Лаг от события до запроса — часы. Двойное хранение: hot data в Kafka + cold data в lake.Analytics Query
Аналитический query (Spark, Trino) видит данные с задержкой в несколько часов. Snapshot-isolation отсутствует — partial reads возможны во время batch-jobs.New: Kafka → Iceberg
Streaming-lakehouse: Iceberg sink connector или Tableflow непрерывно пишет файлы в object storage и атомарно коммитит snapshot. Kafka topic эквивалентен Iceberg table — две проекции одного и того же лога событий.Analytics Query
Запрос видит данные через commit interval (5 минут default для open-source connector, секунды для Tableflow). Snapshot-isolation встроена в Iceberg spec.Apache Iceberg Sink Connector
IcebergSinkConnector — официальный open-source sink, разработанный командой Tabular (acquired Databricks в 2024) и теперь являющийся частью Apache Iceberg core. Доступен в apache/iceberg репозитории, поставляется как plugin для Kafka Connect.
Ключевые свойства:
- Exactly-once семантика через KIP-447 (Kafka 2.5+) — данные пишутся транзакционно: data files + offset commit в одной Kafka-транзакции.
- Multi-table routing — один коннектор может писать в несколько Iceberg-таблиц (fan-out), маршрутизация по полю или по топику.
- Schema evolution — автоматическое добавление новых столбцов на основе схемы из Schema Registry.
- CDC mode — обработка Debezium-событий с upsert/delete вместо append-only.
- REST/Glue/Hive/Nessie/JDBC catalog — поддержка всех основных каталогов.
Минимальная конфигурация
{
"name": "iceberg-events-sink",
"config": {
"connector.class": "org.apache.iceberg.connect.IcebergSinkConnector",
"tasks.max": "4",
"topics": "user-events",
"iceberg.tables": "analytics.user_events",
"iceberg.catalog.type": "rest",
"iceberg.catalog.uri": "https://catalog.example.com",
"iceberg.catalog.warehouse": "s3://my-warehouse/",
"iceberg.catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO",
"iceberg.tables.auto-create-enabled": "true",
"iceberg.tables.evolve-schema-enabled": "true",
"iceberg.control.commit.interval-ms": "300000",
"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false"
}
}
Параметры разбираются по группам:
| Группа | Параметры | Назначение |
|---|---|---|
| Catalog | iceberg.catalog.type, iceberg.catalog.uri, iceberg.catalog.warehouse | Подключение к Iceberg REST/Glue/Hive каталогу |
| Tables | iceberg.tables, iceberg.tables.dynamic-enabled, iceberg.tables.route-field | Статическая или динамическая маршрутизация |
| Schema | iceberg.tables.auto-create-enabled, iceberg.tables.evolve-schema-enabled | Автоматическое создание и эволюция схемы |
| Commit | iceberg.control.commit.interval-ms, iceberg.control.commit.timeout-ms | Частота snapshot-коммитов |
Для production используйте REST catalog (Polaris, Lakekeeper, Snowflake Open Catalog) вместо Hive/JDBC. REST catalog отделяет каталожные операции от storage-операций и позволяет ротировать credentials, не перезапуская коннектор.
Commit strategy и exactly-once
В отличие от обычного sink connector, который коммитит offset после каждого put(), Iceberg sink использует двухэтапный коммит координатора и worker-задач.
Coordinator Task
Coordinator: одна из task в connector выполняет роль координатора. Раз в commit-interval (default 5 минут) рассылает begin-commit event в control topic. Все workers получают событие синхронно.control topic
Control topic: специальный Kafka-топик (default control-iceberg) для координации между coordinator и workers. Coordinator пишет begin-commit, workers отвечают data-files событиями со списком созданных файлов.Worker 1
Worker Task 1: завершает текущий writer (закрывает Parquet файл), отправляет data-files event со списком файлов и Kafka offset, открывает новый writer для следующего commit-периода.Worker 2
Worker Task 2: то же самое — закрывает текущие writers, рапортует список файлов и offset через control topic.Worker N
Worker Task N: каждый worker независимо закрывает свои файлы, но коммит в Iceberg происходит атомарно через координатора.Iceberg snapshot commit
Iceberg snapshot commit: координатор собирает все data-files events, формирует один snapshot со всеми файлами и атомарно коммитит его в Iceberg каталог. Snapshot становится видимым для readers (Spark/Trino) одномоментно.Алгоритм exactly-once:
- Coordinator раз в
commit.interval-msрассылаетbegin-commitevent в control topic. - Каждый worker получает событие, закрывает текущий Parquet writer и отправляет
data-filesevent со списком файлов и Kafka-offset. - Coordinator собирает все data-files events и формирует один Iceberg snapshot атомарно. В тот же момент Kafka offset коммитится транзакционно (KIP-447).
- Если crash случился до коммита — workers удалят оркфаненные файлы при перезапуске и переписали данные. Если после — Iceberg snapshot уже виден читателям, offset закоммичен.
commit.interval-ms (default 300000 = 5 минут) задаёт минимальный лаг видимости данных. Меньшие значения (например, 30 секунд) ускоряют видимость, но создают много мелких файлов — увеличивают overhead на metadata и снижают query performance. Балансируйте: для real-time analytics 1–2 минуты, для daily reports — 10–15 минут.
Schema evolution
Если iceberg.tables.evolve-schema-enabled=true, коннектор автоматически расширяет схему таблицы при появлении новых полей в сообщениях.
| Изменение в сообщении | Поведение коннектора | Iceberg операция |
|---|---|---|
| Добавлено новое поле | Добавляет столбец, старые строки имеют NULL | ALTER TABLE ADD COLUMN |
| Удалено поле | Игнорирует — столбец остаётся в таблице | Не меняет схему |
| Изменён тип (int → long) | Расширяет тип, если совместим | ALTER TABLE ALTER COLUMN TYPE |
| Несовместимое изменение | Ошибка, запись попадает в DLQ | — |
Schema evolution в Iceberg реализована на уровне metadata: добавление столбца не переписывает старые файлы. Reader подставляет NULL для отсутствующих столбцов. Это бесплатная операция, в отличие от классических Hive-таблиц, где требовался полный rewrite.
Multi-table routing и CDC
Один коннектор может писать в несколько Iceberg-таблиц. Два режима:
Static routing — список таблиц в iceberg.tables, маршрутизация по топику:
"topics": "orders,payments,refunds",
"iceberg.tables": "warehouse.orders,warehouse.payments,warehouse.refunds"
Dynamic routing — таблица определяется значением поля в сообщении:
"iceberg.tables.dynamic-enabled": "true",
"iceberg.tables.route-field": "table_name"
При dynamic routing каждое сообщение содержит поле table_name, и коннектор маршрутизирует в соответствующую таблицу. Это полезно для Debezium CDC: одна Debezium-source пишет события из всех таблиц БД в один топик с метаданными, а Iceberg sink раскладывает их по таблицам lakehouse.
CDC mode
Для Debezium-событий включается iceberg.tables.cdc-field:
"iceberg.tables.cdc-field": "op",
"iceberg.tables.upsert-mode-enabled": "true"
Коннектор интерпретирует op (insert / update / delete) и применяет equality deletes + new data files. Iceberg merge-on-read формат сохраняет delete-маркеры, которые применяются при чтении.
Confluent Tableflow
В 2025 году Confluent объявил GA своего управляемого аналога — Tableflow. Концепция: каждый Kafka topic в Confluent Cloud можно одной кнопкой превратить в Iceberg или Delta Lake таблицу.
Tableflow GA включает:
- Iceberg + Delta Lake одновременно — один топик может быть представлен в обоих форматах параллельно.
- Catalog integration — автопубликация в AWS Glue, Snowflake Open Catalog (Polaris), Databricks Unity Catalog.
- Built-in maintenance — автоматическая компакция мелких файлов, snapshot expiration, оркфаненные файлы.
- Dead Letter Queue — schema-violation сообщения изолируются без остановки потока.
- Bring Your Own Key — customer-managed encryption keys для compliance.
Архитектура Tableflow
В отличие от Iceberg sink connector, Tableflow не использует Kafka consumer API. Сервис читает сегменты напрямую из tiered storage Confluent Cloud.
Kafka Topic (Kora)
Kafka topic в Confluent Cloud: использует tiered storage — hot segments на брокерах, cold segments в S3/GCS/Azure Blob. Tableflow знает структуру log-сегментов и читает их напрямую, минуя consumer-API.Tableflow Job
Tableflow materialization job: параллельно читает log-сегменты из tiered storage, применяет схему из Schema Registry, конвертирует в Parquet, формирует Iceberg/Delta data files. Использует Kora cloud-native engine для эластичного масштабирования.Object Storage
Object storage (Confluent-managed или customer-managed): хранит Parquet data files и Iceberg/Delta metadata. Tableflow автоматически компактит мелкие файлы и удаляет старые snapshot.AWS Glue / Polaris / Unity
External catalogs: AWS Glue, Snowflake Open Catalog (Polaris), Databricks Unity Catalog. Tableflow публикует table metadata в выбранный каталог, и таблица сразу доступна Spark/Trino/Snowflake/Databricks.Tableflow читает segment-файлы напрямую из tiered storage, минуя consumer API. Это даёт большую пропускную способность и не нагружает брокеров. Компромисс — это работает только в Confluent Cloud (vendor lock-in).
Schema discovery в Tableflow автоматическая: сервис подтягивает схему из Confluent Schema Registry (Avro/Protobuf/JSON Schema) и маппит её на Iceberg/Delta schema. Schema evolution применяется к таблице при изменении в Schema Registry.
Сравнение подходов
В production выбор сводится к четырём вариантам:
| Решение | Latency | Vendor lock | Maintenance | Cost model |
|---|---|---|---|---|
| Iceberg Sink Connector (Apache) | 5+ минут (commit interval) | Open-source, любой Kafka | Ручная (компакция, snapshot expire) | Compute коннектор-кластера |
| Confluent Tableflow | Секунды | Confluent Cloud only | Полностью managed | Per-GB ingestion + storage |
| Debezium Server + Iceberg sink | 5+ минут | Debezium-only (CDC use case) | Ручная | Compute |
| Kafka → Spark/Flink → Iceberg | Секунды (Flink) / минуты (Spark) | Open-source, гибко | Ручная (Spark/Flink job + табличная) | Compute обоих кластеров |
| Аспект | Iceberg Sink Connector | Tableflow |
|---|---|---|
| Routing | 1:N через route-field, fan-out | 1:1 (один топик — одна таблица) |
| Format | Iceberg | Iceberg + Delta Lake |
| Schema source | Schema Registry или embedded | Schema Registry |
| Catalog | REST/Glue/Hive/Nessie/JDBC | AWS Glue, Polaris, Unity |
| CDC | Через cdc-field + upsert mode | Автоматический CDC stream materialization |
| Maintenance | Внешний (Spark/Trino, Amoro, Nimtable) | Встроенный |
| Минимальный latency | 5 минут (рекомендуется) | Секунды |
Для команды с собственным Kafka-кластером (self-managed) и open-source стеком — Apache Iceberg Sink Connector. Для команды на Confluent Cloud, которая хочет минимальный operational overhead и максимально свежие данные — Tableflow. Для CDC use case — Debezium Server с Iceberg sink (PostgreSQL/MySQL → Iceberg напрямую без Kafka Connect cluster).
Limitations и production tips
Limitations Iceberg Sink Connector
- Commit interval = latency floor. Меньше 1 минуты ставить нельзя без серьёзных проблем с производительностью. 5 минут — здоровый default.
- Append-only по умолчанию. Equality deletes требуют CDC mode и дополнительной нагрузки на reader.
- Maintenance не входит в коннектор. Нужен отдельный job (Spark, Trino, Amoro) для компакции, snapshot expiration, оркфаненных файлов.
- Рост числа small files. При высоком commit-rate каждый snapshot создаёт N файлов (по числу partitions × workers). Без компакции query performance деградирует.
Production tips
{
"iceberg.control.commit.interval-ms": "300000",
"iceberg.tables.default-partition-by": "days(event_time)",
"iceberg.tables.default-properties.write.target-file-size-bytes": "536870912",
"iceberg.tables.default-properties.write.distribution-mode": "hash"
}
| Параметр | Значение | Зачем |
|---|---|---|
commit.interval-ms | 5 минут | Баланс latency и количества файлов |
target-file-size-bytes | 512 MB | Оптимальный размер для Spark/Trino reader |
default-partition-by | days(event_time) | Partition pruning при time-range queries |
write.distribution-mode | hash | Распределение записей по workers по hash-ключу — меньше small files |
Maintenance jobs
Запустите отдельные Spark/Trino задачи на расписании:
-- Компакция мелких файлов
CALL system.rewrite_data_files('analytics.user_events',
options => map('target-file-size-bytes', '536870912'));
-- Удаление старых snapshot (старше 7 дней)
CALL system.expire_snapshots('analytics.user_events',
older_than => TIMESTAMP '2025-04-23 00:00:00');
-- Удаление оркфаненных файлов
CALL system.remove_orphan_files(table => 'analytics.user_events');
Без regular maintenance Iceberg-таблица деградирует за недели: тысячи мелких файлов, гигантская metadata, медленные queries. Включайте maintenance с первого дня — это не “когда-нибудь потом”, а часть production-deployment.
Когда какой подход
| Сценарий | Решение |
|---|---|
| Self-managed Kafka, нужен open-source путь | Apache Iceberg Sink Connector + Polaris/Lakekeeper |
| Confluent Cloud, минимум операций, секунды latency | Tableflow |
| CDC из PostgreSQL/MySQL прямо в Iceberg | Debezium Server + Iceberg sink (без Kafka Connect cluster) |
| Сложные трансформации между Kafka и Iceberg | Flink Dynamic Iceberg Sink (joins, aggregates, filtering) |
| Daily batch enough, hadoop-стек уже есть | Spark Structured Streaming + Iceberg writer |
Ключевые выводы
- Streaming-Lakehouse pattern в 2025 — стандарт: Kafka topic эквивалентен Iceberg table через коннектор или managed-сервис.
- Apache Iceberg Sink Connector — open-source путь, exactly-once через KIP-447, commit-interval (default 5 минут) задаёт минимальный latency.
- Confluent Tableflow — управляемый сервис в Confluent Cloud, читает segment-файлы напрямую, секунды latency, встроенный maintenance, Iceberg + Delta Lake одновременно.
- Schema evolution автоматическая в обоих подходах: новые столбцы добавляются на metadata-уровне, без переписывания файлов.
- CDC mode в Iceberg sink через
cdc-field+ upsert + equality deletes; в Tableflow CDC поддерживается из коробки. - Production maintenance обязателен: компакция файлов, snapshot expiration, orphan files. Tableflow делает это сам, для open-source connector нужен отдельный job.
- REST catalog (Polaris, Lakekeeper, Snowflake Open Catalog, Unity) — стандарт для streaming-lakehouse в 2025; статические Hive/JDBC каталоги уходят.