skip to content

How does Kafka Connect track offsets for source connectors versus sink connectors? Where is each kind stored?

level: juniorimportance: must knowfreq 75%

answer

  1. source = connector-defined offsets -> offset.storage.topic
  2. sink = it's a consumer -> __consumer_offsets
  3. consumer group connect-<name>
  4. sourcePartition + sourceOffset maps
  5. OffsetStorageReader on resume

basics

~20 s

Source connectors store their own progress (where they read from the external system) in a special Kafka topic called the offset storage topic. Sink connectors just use a normal Kafka consumer group, so Kafka tracks their offsets like any consumer.

solid answer

~40 s

Source and sink connectors track progress differently because they move data in opposite directions. A source connector reads from an external system (a database, a file, an API) and writes to Kafka; Connect has no built-in notion of 'position' in that external system, so the connector emits an arbitrary source partition + source offset map (e.g. {filename} -> {byte position}) that the framework persists to the offset.storage.topic (a compacted Kafka topic). A sink connector reads from Kafka and writes outward, so it IS a Kafka consumer: it joins a consumer group named connect-<connector-name> and its progress is the normal committed consumer-group offsets stored in __consumer_offsets. So source offsets are connector-defined and live in offset.storage.topic; sink offsets are Kafka consumer-group offsets.

go deeper

for a junior

Know the one-line split: source offsets in a special Connect topic, sink offsets in a normal consumer group.

for a middle

Explain sourcePartition/sourceOffset opaque maps and the connect-<name> consumer group; know the standalone file variant.

for a senior

Tie storage location to how you reset/inspect each (REST offsets API vs kafka-consumer-groups), and the compacted nature of offset.storage.topic.

for a principal

Reason about the framework contract (OffsetStorageReader on resume, preCommit on sinks) and why the asymmetry is fundamental to direction of data flow.

## The core asymmetry Kafka Connect runs two kinds of connectors, and they track 'where am I' completely differently because data flows in opposite directions. **Source connectors** pull data FROM an external system (a database, a file, an S3 bucket, an API) and produce it INTO Kafka topics. Kafka itself has no idea what 'position' means in your external system. So the connector author defines two opaque maps for every record-producing unit: - a **source partition** — identifies *what* is being read, e.g. `{"filename": "app.log"}` or `{"table": "orders"}`. - a **source offset** — identifies *how far* you've read within that partition, e.g. `{"position": 10342}` (a byte offset) or `{"timestamp": 1719... , "id": 5567}` (a DB cursor). The Connect framework persists these maps for you. In **distributed mode** they go to a compacted Kafka topic named by `offset.storage.topic`; in **standalone mode** they go to a local file named by `offset.storage.file.filename`. On restart, the framework hands the last committed source offset back to the task (via `SourceTaskContext.offsetStorageReader()`) so it can resume from the right place instead of re-reading everything. **Sink connectors** do the opposite: they CONSUME from Kafka topics and write OUT to an external system. A sink task *is literally a Kafka consumer*. The framework subscribes it to the configured topics under a consumer group named `connect-<connector-name>`. 'How far have we processed' is therefore just the **committed consumer-group offsets**, stored in Kafka's internal `__consumer_offsets` topic — exactly like any other consumer application. You can inspect them with `kafka-consumer-groups.sh --group connect-my-sink --describe`. ## Why this matters operationally - To reset a **source** connector you manipulate `offset.storage.topic` (or use the Connect REST offsets API in modern versions). - To reset a **sink** connector you manipulate its consumer group offsets (same API in modern Connect, or `kafka-consumer-groups.sh --reset-offsets` when stopped). - The two storage locations are independent; a sink connector does not write to `offset.storage.topic` at all, and a source connector does not appear as a consumer group. ## Edge case Sink connectors *can* additionally read source-style offsets via `SinkTaskContext` only in special cases; normally the consumer group is the single source of truth, and Connect commits those offsets on the `offset.flush.interval.ms` cadence (or per-record/SMT-driven `preCommit`).

  • What is the consumer group name for a sink connector called 'es-sink'?
    connect-es-sink — Connect names the group connect-<connector-name>, and you can describe/reset it with kafka-consumer-groups.sh just like any consumer group.
  • In standalone mode, where do source offsets go instead of the topic?
    To a local file set by offset.storage.file.filename; there is no offset.storage.topic in standalone mode.

saying these in an interview costs you the question

  • Saying both source and sink offsets live in offset.storage.topic (sinks use __consumer_offsets).
  • Saying sink connectors define their own offset maps (they're plain consumer-group offsets).
  • Claiming Connect understands the external system's position automatically (the connector author defines the source partition/offset).

context