Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 06.06 · 25 мин
Продвинутый
Apache IcebergIceberg Sink ConnectorConfluent TableflowStreaming LakehouseOpen Table FormatDelta LakeREST Catalog

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, как только пакет готов.

NOTE

Apache Iceberg по своей природе таблица: каждая строка имеет первичный ключ, schema-evolution встроена в спецификацию, snapshot изолированы. Kafka — упорядоченный лог. Соединение этих абстракций (Kafka topic как Iceberg table) — главный архитектурный паттерн 2025 года, который Jack Vanlightly называет “stream/table duality”.

Streaming-Lakehouse: до и после

Old: Kafka + Batch ETL

Классический подход: Kafka хранит события 7 дней, batch ETL раз в час читает их в Spark, дедуплицирует, пишет в Parquet/Hive. Лаг от события до запроса — часы. Двойное хранение: hot data в Kafka + cold data в lake.
hours latency

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 — две проекции одного и того же лога событий.
seconds–minutes

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"
  }
}

Параметры разбираются по группам:

ГруппаПараметрыНазначение
Catalogiceberg.catalog.type, iceberg.catalog.uri, iceberg.catalog.warehouseПодключение к Iceberg REST/Glue/Hive каталогу
Tablesiceberg.tables, iceberg.tables.dynamic-enabled, iceberg.tables.route-fieldСтатическая или динамическая маршрутизация
Schemaiceberg.tables.auto-create-enabled, iceberg.tables.evolve-schema-enabledАвтоматическое создание и эволюция схемы
Commiticeberg.control.commit.interval-ms, iceberg.control.commit.timeout-msЧастота snapshot-коммитов
TIP

Для 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-задач.

Iceberg Sink commit coordination

Coordinator Task

Coordinator: одна из task в connector выполняет роль координатора. Раз в commit-interval (default 5 минут) рассылает begin-commit event в control topic. Все workers получают событие синхронно.
begin-commit event

control topic

Control topic: специальный Kafka-топик (default control-iceberg) для координации между coordinator и workers. Coordinator пишет begin-commit, workers отвечают data-files событиями со списком созданных файлов.
workers receive event

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 происходит атомарно через координатора.
data-files events

Iceberg snapshot commit

Iceberg snapshot commit: координатор собирает все data-files events, формирует один snapshot со всеми файлами и атомарно коммитит его в Iceberg каталог. Snapshot становится видимым для readers (Spark/Trino) одномоментно.

Алгоритм exactly-once:

  1. Coordinator раз в commit.interval-ms рассылает begin-commit event в control topic.
  2. Каждый worker получает событие, закрывает текущий Parquet writer и отправляет data-files event со списком файлов и Kafka-offset.
  3. Coordinator собирает все data-files events и формирует один Iceberg snapshot атомарно. В тот же момент Kafka offset коммитится транзакционно (KIP-447).
  4. Если crash случился до коммита — workers удалят оркфаненные файлы при перезапуске и переписали данные. Если после — Iceberg snapshot уже виден читателям, offset закоммичен.
WARNING

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 операция
Добавлено новое полеДобавляет столбец, старые строки имеют NULLALTER TABLE ADD COLUMN
Удалено полеИгнорирует — столбец остаётся в таблицеНе меняет схему
Изменён тип (int → long)Расширяет тип, если совместимALTER TABLE ALTER COLUMN TYPE
Несовместимое изменениеОшибка, запись попадает в DLQ
NOTE

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.

Tableflow internal architecture

Kafka Topic (Kora)

Kafka topic в Confluent Cloud: использует tiered storage — hot segments на брокерах, cold segments в S3/GCS/Azure Blob. Tableflow знает структуру log-сегментов и читает их напрямую, минуя consumer-API.
read segment files

Tableflow Job

Tableflow materialization job: параллельно читает log-сегменты из tiered storage, применяет схему из Schema Registry, конвертирует в Parquet, формирует Iceberg/Delta data files. Использует Kora cloud-native engine для эластичного масштабирования.
write Parquet

Object Storage

Object storage (Confluent-managed или customer-managed): хранит Parquet data files и Iceberg/Delta metadata. Tableflow автоматически компактит мелкие файлы и удаляет старые snapshot.
publish metadata

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).

INFO

Schema discovery в Tableflow автоматическая: сервис подтягивает схему из Confluent Schema Registry (Avro/Protobuf/JSON Schema) и маппит её на Iceberg/Delta schema. Schema evolution применяется к таблице при изменении в Schema Registry.


Сравнение подходов

В production выбор сводится к четырём вариантам:

РешениеLatencyVendor lockMaintenanceCost model
Iceberg Sink Connector (Apache)5+ минут (commit interval)Open-source, любой KafkaРучная (компакция, snapshot expire)Compute коннектор-кластера
Confluent TableflowСекундыConfluent Cloud onlyПолностью managedPer-GB ingestion + storage
Debezium Server + Iceberg sink5+ минутDebezium-only (CDC use case)РучнаяCompute
Kafka → Spark/Flink → IcebergСекунды (Flink) / минуты (Spark)Open-source, гибкоРучная (Spark/Flink job + табличная)Compute обоих кластеров
АспектIceberg Sink ConnectorTableflow
Routing1:N через route-field, fan-out1:1 (один топик — одна таблица)
FormatIcebergIceberg + Delta Lake
Schema sourceSchema Registry или embeddedSchema Registry
CatalogREST/Glue/Hive/Nessie/JDBCAWS Glue, Polaris, Unity
CDCЧерез cdc-field + upsert modeАвтоматический CDC stream materialization
MaintenanceВнешний (Spark/Trino, Amoro, Nimtable)Встроенный
Минимальный latency5 минут (рекомендуется)Секунды
TIP

Для команды с собственным 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

  1. Commit interval = latency floor. Меньше 1 минуты ставить нельзя без серьёзных проблем с производительностью. 5 минут — здоровый default.
  2. Append-only по умолчанию. Equality deletes требуют CDC mode и дополнительной нагрузки на reader.
  3. Maintenance не входит в коннектор. Нужен отдельный job (Spark, Trino, Amoro) для компакции, snapshot expiration, оркфаненных файлов.
  4. Рост числа 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-ms5 минутБаланс latency и количества файлов
target-file-size-bytes512 MBОптимальный размер для Spark/Trino reader
default-partition-bydays(event_time)Partition pruning при time-range queries
write.distribution-modehashРаспределение записей по 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');
WARNING

Без regular maintenance Iceberg-таблица деградирует за недели: тысячи мелких файлов, гигантская metadata, медленные queries. Включайте maintenance с первого дня — это не “когда-нибудь потом”, а часть production-deployment.


Когда какой подход

СценарийРешение
Self-managed Kafka, нужен open-source путьApache Iceberg Sink Connector + Polaris/Lakekeeper
Confluent Cloud, минимум операций, секунды latencyTableflow
CDC из PostgreSQL/MySQL прямо в IcebergDebezium Server + Iceberg sink (без Kafka Connect cluster)
Сложные трансформации между Kafka и IcebergFlink Dynamic Iceberg Sink (joins, aggregates, filtering)
Daily batch enough, hadoop-стек уже естьSpark Structured Streaming + Iceberg writer

Ключевые выводы

  1. Streaming-Lakehouse pattern в 2025 — стандарт: Kafka topic эквивалентен Iceberg table через коннектор или managed-сервис.
  2. Apache Iceberg Sink Connector — open-source путь, exactly-once через KIP-447, commit-interval (default 5 минут) задаёт минимальный latency.
  3. Confluent Tableflow — управляемый сервис в Confluent Cloud, читает segment-файлы напрямую, секунды latency, встроенный maintenance, Iceberg + Delta Lake одновременно.
  4. Schema evolution автоматическая в обоих подходах: новые столбцы добавляются на metadata-уровне, без переписывания файлов.
  5. CDC mode в Iceberg sink через cdc-field + upsert + equality deletes; в Tableflow CDC поддерживается из коробки.
  6. Production maintenance обязателен: компакция файлов, snapshot expiration, orphan files. Tableflow делает это сам, для open-source connector нужен отдельный job.
  7. REST catalog (Polaris, Lakekeeper, Snowflake Open Catalog, Unity) — стандарт для streaming-lakehouse в 2025; статические Hive/JDBC каталоги уходят.
Проверка знанийKnowledge check
Команда настроила Apache Iceberg Sink Connector с iceberg.control.commit.interval-ms=30000 (30 секунд) для real-time дашбордов. Через две недели query performance в Trino упала с секунд до десятков секунд, хотя данных не сильно прибавилось. В чём причина и как чинить?
ОтветAnswer
Причина — small files problem. Каждый commit создаёт snapshot с N data files (по числу partitions × workers). При commit interval 30 секунд за две недели создаются десятки тысяч мелких Parquet-файлов. Каждый файл — отдельный read, отдельный footer-parse, отдельный network roundtrip. Trino reader тратит больше времени на metadata и open-file overhead, чем на чтение данных. Решение в три шага: (1) увеличить commit interval до 2–5 минут — баланс latency vs file count; (2) запустить Spark/Trino procedure rewrite_data_files регулярно (раз в час или сутки) для компакции в файлы 512 MB; (3) запустить expire_snapshots для удаления старых snapshot и связанных metadata. Альтернатива — мигрировать на Tableflow, который делает компакцию автоматически. Урок: в streaming-lakehouse latency requirement 30 секунд требует built-in maintenance, иначе деградация неизбежна.
Проверка знанийKnowledge check
В чём принципиальное архитектурное отличие Confluent Tableflow от Apache Iceberg Sink Connector в способе чтения данных из Kafka, и почему это даёт Tableflow преимущество в latency?
ОтветAnswer
Apache Iceberg Sink Connector использует стандартный Kafka consumer API: подписывается на топик через consumer group, получает записи через poll(), сериализует, пишет в Iceberg. Каждый шаг — стандартный Kafka-протокол с накладными расходами на сетевые roundtrip, deserialization, consumer rebalance. Tableflow читает log-сегменты напрямую из tiered storage Confluent Cloud (S3/GCS/Azure), минуя consumer API. Сервис знает внутреннюю структуру log-сегментов и параллельно читает их как файлы из object storage. Это даёт два преимущества: (1) пропускная способность ограничена только object-storage bandwidth, а не консьюмером; (2) брокера Kafka не нагружаются — Tableflow не конкурирует за ресурсы с production-консьюмерами. Компромисс — vendor lock-in: эта оптимизация работает только в Confluent Cloud, где сервис имеет привилегированный доступ к tiered storage и log-формату Kora.
Проверка знанийKnowledge check
Команда переходит с PostgreSQL на streaming-lakehouse: события из БД должны попадать в Iceberg-таблицу с поддержкой UPDATE и DELETE (не только INSERT), потому что строки в источнике обновляются. Какие два настройки Iceberg Sink Connector обязательно нужны и какой формат Iceberg-таблицы это создаёт?
ОтветAnswer
Нужны iceberg.tables.cdc-field=op (или другое поле) и iceberg.tables.upsert-mode-enabled=true. cdc-field указывает поле сообщения, содержащее тип операции (insert / update / delete) — обычно из Debezium envelope. upsert-mode-enabled включает merge-on-read семантику: коннектор для UPDATE генерирует equality delete (по primary key) + новый data file, для DELETE — только equality delete. Таблица становится merge-on-read форматом: data files (новые или обновлённые строки) + delete files (маркеры удалений). Reader (Spark/Trino) при чтении применяет deletes к data files на лету. Это медленнее append-only, но позволяет CDC pattern. Production-tip: запускайте rewrite_data_files регулярно с опцией delete-file-threshold — это переписывает файлы и применяет deletes физически, восстанавливая read performance. Без compaction merge-on-read деградирует быстро при большом числе апдейтов.
Iceberg catalog: REST, Polaris, Lakekeeper

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 7