Explain the pattern of using Kafka as the streaming backbone for ELT into a lakehouse, and how Iceberg/lake table formats fit into it.
answer
- ELT = load raw first, transform in destination
- Kafka = decoupled event highway; CDC + sinks load bronze
- Iceberg = metadata over Parquet: ACID, time travel, evolution
- Kafka→Iceberg sink commits batches as snapshots, upsert for CDC
- small-file compaction
basics
~20 sKafka acts as the central event highway: producers and CDC feed raw events in, sink connectors land them in cheap lake storage, and transformations run afterward in the warehouse/lake (the 'T' of ELT). Table formats like Apache Iceberg make those landed files behave like real, queryable, transactional tables.
solid answer
~50 sIn **ELT** you Extract and Load raw data first, then Transform it inside the destination — the opposite ordering of ETL. Kafka serves as the **streaming backbone**: producers, Debezium CDC, and other sources publish raw events to topics; **sink connectors** continuously load them into the lake (S3/Iceberg/BigQuery/Snowflake) as the raw/bronze layer; then transformations (dbt, Spark, warehouse SQL) build silver/gold tables. This decouples producers from consumers, gives replay/backfill, and turns the lake into the analytics store. **Apache Iceberg** (and Delta/Hudi) is an open **table format**: a metadata layer over Parquet files giving ACID snapshots, schema/partition evolution, time travel, and hidden partitioning — so the landed files behave like a real transactional table that many engines (Spark, Trino, Flink, Snowflake) can read consistently. Modern Kafka→Iceberg sinks (Confluent's Iceberg sink, Tabler/Iceberg Kafka Connect, Flink) commit batches as Iceberg snapshots with upsert/merge for CDC, handling small-file compaction.
go deeper
Know Kafka feeds raw events into the lake and transformations run afterward (ELT); Iceberg makes files into real tables.
Explain the bronze/silver/gold flow, producer-consumer decoupling, and what Iceberg adds over raw Parquet.
Detail Kafka→Iceberg snapshot commits, CDC upserts, compaction, and replay/backfill benefits.
Govern table-format choice, multi-engine consistency, lock-in, and lakehouse layering at org scale.
## ETL vs ELT - **ETL** (Extract, Transform, Load): clean/shape data *before* loading into the warehouse. Transformation logic lives in a separate engine; the warehouse only stores final tables. - **ELT** (Extract, **Load**, **Transform**): load raw data first, then transform *inside* the destination using its compute (cheap cloud storage + scalable query engines make this practical). You keep the raw data, transformations are versioned SQL (e.g., dbt), and you can re-transform anytime. The modern cloud lakehouse strongly favors ELT. ## Kafka as the backbone Kafka sits at the center as the **event highway**: 1. **Extract/produce**: application producers, Debezium **CDC** from operational DBs, and other sources publish raw events to **topics**. Kafka decouples producers from consumers — many independent sinks/consumers can read the same stream. 2. **Load**: **Kafka Connect sink connectors** continuously stream those events into the lake as the **raw / bronze layer** (e.g., Parquet in S3, or directly into Iceberg/BigQuery/Snowflake). This is the 'L'. 3. **Transform**: downstream tools (dbt, Spark, Flink, warehouse SQL) build cleaned **silver** and aggregated **gold** tables. This is the 'T', running *after* load, in the destination. Because Kafka retains the stream (and with **tiered storage**, cheaply for long periods), you get **replay** and **backfill**: rebuild a transformed table from scratch, onboard a new consumer from history, or recover from a bad transform. ## Where table formats fit Raw Parquet files in S3 are just files — no transactions, no consistent schema view, painful partition management. An open **table format** fixes this: - **Apache Iceberg** (also **Delta Lake**, **Apache Hudi**) is a **metadata layer over data files** (usually Parquet). It provides: - **ACID snapshots**: each commit is an atomic snapshot; readers see consistent data, concurrent writers don't corrupt each other. - **Schema and partition evolution**: add/drop/rename columns and change partitioning without rewriting all data. - **Time travel**: query the table as of an earlier snapshot/timestamp. - **Hidden partitioning**: queries prune partitions without users hardcoding partition columns. - **Engine-agnostic**: Spark, Trino, Flink, Snowflake, BigQuery (via external/managed Iceberg) all read the same table. This turns the streamed-in files into a **real, queryable, transactional table** — the heart of the lakehouse. ## Kafka → Iceberg specifics Dedicated sinks (the **Confluent Iceberg sink**, the community **Iceberg Kafka Connect** connector, or **Flink** writing Iceberg) consume topics and **commit record batches as Iceberg snapshots**. For CDC streams they perform **upsert/merge** keyed on the primary key (applying inserts/updates/deletes) so the Iceberg table reflects current state, and they coordinate **small-file compaction** (streaming writes naturally produce many tiny files that hurt query performance, so compaction rewrites them into larger files). ## Why this pattern wins - One durable, replayable backbone feeds *all* downstream consumers consistently. - Raw retained + ELT means transformations are reproducible and auditable. - Open table formats avoid vendor lock-in and let multiple engines share one source of truth.
- Why is a table format like Iceberg needed on top of raw Parquet files in S3?Raw Parquet files have no transactions, no atomic multi-file commits, no consistent schema/partition view, and no time travel. Iceberg adds an ACID metadata layer giving snapshot isolation, schema/partition evolution, and engine-agnostic consistent reads — making the files behave like a real transactional table.
- What operational problem do streaming writes into Iceberg create, and how is it handled?Continuous streaming produces many small files, which slows queries (lots of metadata and file opens). It's handled by compaction jobs that rewrite small files into larger ones and expire old snapshots, often run on a schedule alongside the sink.
saying these in an interview costs you the question
- Confusing ELT with ETL — in ELT the transform happens after loading, in the destination.
- Calling Iceberg a file format — it's a table/metadata format layered over Parquet files.
- Ignoring small-file compaction for streaming writes.
- Saying raw S3 Parquet alone gives ACID/time-travel — it doesn't without a table format.