skip to content

Offset Management in Connect

Where Connect keeps its progress: source offsets in an internal topic, sink offsets in a consumer group, and how to reset either. Comes up whenever a pipeline has to be replayed from a point in time.

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

questions

5

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

open as a page

Walk through how a source task's offsets get committed. What does offset.flush.interval.ms control, and what is the relationship between producing records and committing offsets?

level: middleimportance: must knowfreq 60%

basics

~20 s

A source task hands records to the framework, which produces them to Kafka. Periodically (every offset.flush.interval.ms, default 60000 ms) the framework flushes the producer, then writes the source offsets of acknowledged records to offset.storage.topic. Offsets are only committed after the data is safely in Kafka.

open as a page

What is the OffsetStorageReader, and how should a source task use it on startup?

level: middleimportance: should knowfreq 45%

basics

~20 s

OffsetStorageReader lets a source task read back the last committed source offsets when it starts, so it can resume from where it left off. The task gets it from its SourceTaskContext and queries by source partition.

open as a page

How does exactly-once delivery for SOURCE connectors work in Kafka Connect, and what is required to enable it?

level: seniorimportance: should knowfreq 40%

basics

~20 s

Connect (KIP-618, Kafka 3.3+) can run source connectors exactly-once by writing the produced records and their source offsets together in a single Kafka transaction. You enable it cluster-wide with exactly.once.source.support=enabled, and the connector must declare exactly.once.support and define transaction boundaries.

open as a page

A team needs to make a source connector reprocess data from the beginning, and separately rewind a sink connector. How do you reset offsets for each, safely?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Stop the connector first. For modern Connect (3.6+), use the REST offsets API: GET/DELETE/PATCH /connectors/{name}/offsets. DELETE wipes offsets so a source restarts from scratch; PATCH sets specific source offsets or sink consumer-group offsets. Older clusters edit offset.storage.topic (source) or use kafka-consumer-groups --reset-offsets (sink).

open as a page