skip to content

Kafka in the Data Ecosystem

Kafka as the pipe into a lakehouse or warehouse: Connect sinks to S3, Iceberg or Snowflake, CDC pipelines, and tiered storage. Interviewers ask how streaming data ends up queryable by analysts.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

6

What is Kafka Connect, and how do sink connectors get data from Kafka topics into a data lake or warehouse like S3, BigQuery, or Snowflake?

level: juniorimportance: must knowfreq 70%

answer

  1. framework not custom code
  2. source in / sink out
  3. tasks.max = parallelism
  4. converter + Schema Registry
  5. at-least-once by default

basics

~20 s

Kafka Connect is a framework for moving data between Kafka and external systems without custom code. A sink connector reads records from Kafka topics and writes them to a destination like S3, BigQuery, or Snowflake using ready-made plugins and config.

solid answer

~40 s

Kafka Connect is a runtime and framework for streaming data into and out of Kafka declaratively via JSON config instead of bespoke producer/consumer code. A source connector pulls data into Kafka; a sink connector reads from topics and writes to an external system. For lakehouse ingestion you run sink connectors such as the Confluent S3 sink, BigQuery sink, or Snowflake sink. Connect runs in distributed mode as a cluster of workers; each connector is split into tasks for parallelism, and offsets are committed so delivery is at-least-once by default. Sinks buffer records and flush in batches (e.g., S3 sink writes files partitioned by topic/partition/time). Converters (Avro/JSON/Protobuf, often with Schema Registry) handle serialization, and Single Message Transforms (SMTs) can reshape records in flight.

go deeper

for a junior

Know Connect = no-code framework, sink = topic to external system, and name S3/BigQuery/Snowflake sinks.

for a middle

Explain workers, tasks.max parallelism, converters + Schema Registry, and the at-least-once flush/commit cycle.

for a senior

Discuss partitioners, file formats, DLQ, SMTs, and how at-least-once duplicates are handled downstream.

for a principal

Reason about Connect cluster sizing, schema-evolution governance, and when a managed sink beats a bespoke streaming job.

## What problem Kafka Connect solves Writing a one-off consumer to copy every Kafka topic into a warehouse is repetitive, error-prone work: you re-solve batching, retries, offset tracking, schema handling, and scaling each time. **Kafka Connect** is a framework (shipped with Apache Kafka) that turns this into configuration. You declare *what* to connect and Connect handles *how*. ## Source vs sink - A **source connector** brings external data *into* Kafka topics (e.g., a database into Kafka). - A **sink connector** reads *from* Kafka topics and writes data *out* to an external system (S3, BigQuery, Snowflake, Iceberg, Elasticsearch, JDBC databases). For lakehouse/warehouse ingestion you almost always use **sink connectors**. ## How a sink works mechanically 1. **Workers**: Connect runs as a cluster of JVM processes called *workers*. In **distributed mode** they coordinate via internal Kafka topics that store connector configs, offsets, and status, so the cluster is fault-tolerant and rebalances if a worker dies. 2. **Connector → tasks**: A connector is configured once but is split into N **tasks** (`tasks.max`) for parallelism. Connect assigns topic partitions across tasks much like a consumer group, so throughput scales with partitions. 3. **Converters**: Bytes on the wire must be deserialized. A **converter** (e.g., `AvroConverter`, `JsonConverter`, `ProtobufConverter`) turns raw bytes into Connect's internal record format. With a **Schema Registry**, Avro/Protobuf converters validate and evolve schemas. 4. **SMTs (Single Message Transforms)**: Lightweight per-record transforms (rename fields, add a timestamp, route topics, mask data) applied before the record reaches the sink. 5. **Buffering and flush**: Sinks accumulate records and **flush in batches**. The S3 sink, for example, groups records by topic/partition and a partitioner (default by Kafka partition, or `TimeBasedPartitioner` for date paths like `year=/month=/day=/hour=`) and writes objects when `flush.size` records accumulate or a `rotate.interval.ms` elapses, in formats like Parquet, Avro, or JSON. 6. **Offset commit**: After a successful flush, Connect commits consumer offsets. Because the data is written *then* offsets committed, the default guarantee is **at-least-once**: a crash between write and commit can replay records, producing duplicates downstream. ## Why this matters for analytics Sink connectors give you a managed, scalable, schema-aware bridge from the streaming layer to the analytics layer with retries, dead-letter queues (`errors.deadletterqueue.topic.name`), and monitoring built in — no custom ingestion service to maintain.

  • What delivery guarantee does a Connect sink give by default, and why?
    At-least-once: the sink writes the batch to the destination, then commits offsets. A failure between those steps replays records, so downstream can see duplicates unless the sink supports exactly-once (e.g., via idempotent upserts or transactional writes).
  • What is the role of a converter versus an SMT?
    A converter (de)serializes the whole record key/value to/from bytes, usually tied to a wire format and Schema Registry. An SMT is a small per-record transformation (rename, mask, route, add field) applied after deserialization, before the connector hands the record to the sink.

saying these in an interview costs you the question

  • Saying Connect requires writing custom producer/consumer code — its whole point is config-driven, no code.
  • Claiming sinks are exactly-once by default — they are at-least-once unless the connector adds dedup/transactional support.
  • Confusing source (into Kafka) and sink (out of Kafka) direction.

context

open as a page

How does a Debezium CDC pipeline move database changes through Kafka into an analytics lakehouse, and what makes it preferable to periodic batch extracts?

level: seniorimportance: must knowfreq 60%

basics

~20 s

Debezium is a Kafka Connect source connector that reads a database's transaction log and emits each insert/update/delete as a Kafka event. Sink connectors then land those change events in the lake, giving near-real-time, low-load, complete change history instead of nightly full-table dumps.

open as a page

Explain the pattern of using Kafka as the streaming backbone for ELT into a lakehouse, and how Iceberg/lake table formats fit into it.

level: middleimportance: should knowfreq 45%

basics

~20 s

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

open as a page

How does the Confluent S3 sink connector lay out files in object storage, and how can it achieve exactly-once delivery to S3?

level: middleimportance: should knowfreq 50%

basics

~20 s

The S3 sink groups records by partitioner (by Kafka partition or by time) and writes batched objects (Parquet/Avro/JSON) when flush.size or a rotate interval is hit. It achieves exactly-once by naming files deterministically from offsets, so replays overwrite the same object instead of duplicating.

open as a page

What is Kafka Tiered Storage (KIP-405), and how does it change the economics and architecture of using Kafka as a long-retention backbone for analytics?

level: seniorimportance: should knowfreq 40%

basics

~20 s

Tiered Storage (KIP-405) lets a Kafka broker offload older log segments to cheap remote object storage (like S3) while keeping recent data on local disk. This makes long retention affordable and lets brokers store far more history without huge local disks.

open as a page

What delivery and schema semantics should you expect from the Snowflake and BigQuery Kafka Connect sinks, and how do they differ from a file-based S3 sink?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

Warehouse sinks (Snowflake, BigQuery) write rows directly into managed tables rather than files in a bucket. They use the warehouse's streaming-ingest APIs, can offer exactly-once or near-exactly-once via offset tracking, and auto-map record schemas to table columns — whereas the S3 sink just writes batched files you must catalog separately.

open as a page