What is the difference between a Source connector and a Sink connector in Kafka Connect, and which direction does data flow in each?
answer
- Source = into Kafka
- Sink = out of Kafka
- SourceRecord vs SinkRecord
- Sink = managed consumer, Source = managed producer
- Kafka is the river
basics
~20 sA Source connector pulls data from an external system INTO Kafka topics. A Sink connector reads data FROM Kafka topics and writes it OUT to an external system. Source = into Kafka, Sink = out of Kafka.
solid answer
~40 sKafka Connect moves data between Kafka and external systems. A Source connector (e.g. Debezium, JDBC source) ingests data from an external system — a database, file, queue — and produces it into Kafka topics; it emits SourceRecords. A Sink connector consumes records from Kafka topics and delivers them to an external sink — S3, Elasticsearch, a JDBC table; it receives SinkRecords. The mnemonic: think of Kafka as the river; the source feeds water in, the sink drains it out. Both run inside Connect worker JVMs, are configured via JSON over the REST API, and split work into Tasks for parallelism. A Source connector's task implements poll() to fetch new data; a Sink connector's task implements put() to receive batches of records.
go deeper
Memorize the direction: source in, sink out, and the two record types.
Know that connectors coordinate while tasks move data, and the poll()/put() split.
Explain the consumer-group (sink) vs producer+offset-topic (source) plumbing differences.
Reason about pipeline topology (DB->Kafka->ES), delivery guarantees, and when to compose source+sink vs a single streaming app.
## What Kafka Connect is Kafka Connect is a framework and runtime for streaming data between Apache Kafka and other systems without writing bespoke producer/consumer code. You deploy *connectors* (plugins) into Connect *workers* (JVM processes) and configure them with JSON via a REST API. ## The two directions There are exactly two kinds of connector, defined by data direction relative to Kafka: - **Source connector** — reads from an *external system* (the source of truth) and **writes into Kafka topics**. Examples: Debezium (change-data-capture from databases), the Confluent JDBC Source, FileStreamSource. The unit it emits is a `SourceRecord`. - **Sink connector** — reads from *Kafka topics* and **writes out to an external system** (the destination). Examples: S3 Sink, Elasticsearch Sink, JDBC Sink. The unit it receives is a `SinkRecord`. Kafka sits in the middle. A useful mental model: Kafka is a river. A **source** pours water *in*; a **sink** drains water *out*. ## How each does its job The top-level class (`SourceConnector` / `SinkConnector`) does NOT move data itself — it is a coordinator. It validates config and produces a list of task configs. The actual data movement is done by **tasks**: - A `SourceTask` implements `poll()`, returning a `List<SourceRecord>` that Connect produces to Kafka. - A `SinkTask` implements `put(Collection<SinkRecord>)`, receiving batches consumed from Kafka that it must deliver downstream. ## Why the distinction matters The two share the framework but differ in their guarantees and plumbing. Sink connectors are essentially managed Kafka *consumers*: Connect runs a consumer group under the hood, so offset tracking, rebalancing, and `topics`/`topics.regex` subscription are handled by the consumer protocol. Source connectors are managed *producers*: Connect stores source-system offsets (e.g. a database log position) in an internal offsets topic so the connector can resume after a restart. This is why only source connectors deal with `taskConfigs()` deciding how to split an external partition space, while sink connectors lean on consumer-group partition assignment. ## Edge cases - A single Connect cluster can run many source and sink connectors simultaneously. - The same topic can be the target of a source connector and the input of a sink connector, forming a pipeline (DB -> Kafka -> Elasticsearch). - Single Message Transforms (SMTs) and converters apply to both directions but in mirror order.
- Which record type does each emit/receive?A SourceTask emits SourceRecords (produced into Kafka); a SinkTask receives SinkRecords (consumed from Kafka).
- Under the hood, what Kafka client does a sink connector use?A sink connector runs as a Kafka consumer group — Connect manages the consumers, so offset commits and partition assignment come from the consumer protocol.
saying these in an interview costs you the question
- Saying a Source connector reads from Kafka — it writes TO Kafka.
- Claiming the Connector class itself moves the data (it only coordinates; Tasks move data).
- Confusing the direction: 'sink ingests into Kafka' is wrong.