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?
answer
- STOP first (PUT /stop), then resume after
- REST: GET/DELETE/PATCH /connectors/{name}/offsets (KIP-875, 3.6+)
- DELETE => source from scratch / sink per auto.offset.reset
- sink PATCH = kafka_topic/kafka_partition/kafka_offset
- legacy sink: kafka-consumer-groups --reset-offsets on connect-<name>
basics
~20 sStop 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).
solid answer
~40 sFirst STOP the connector (PUT /connectors/{name}/stop) — resetting a running connector is unsafe because tasks may re-commit. Then in Kafka 3.6+ use the per-connector offsets REST API: GET /connectors/{name}/offsets to inspect; DELETE /connectors/{name}/offsets to clear all offsets (a source then reprocesses from its initial position; a sink rewinds to its auto.offset.reset); PATCH /connectors/{name}/offsets with a body of offset entries to set exact positions. For a source the entry is {partition:{...sourcePartition}, offset:{...sourceOffset}}; for a sink it's {partition:{kafka_topic, kafka_partition}, offset:{kafka_offset}} since sink offsets ARE consumer-group offsets. Resume with PUT .../resume. On pre-3.6 clusters: for sources, produce tombstones/new values to offset.storage.topic; for sinks, run kafka-consumer-groups.sh --group connect-<name> --reset-offsets --to-earliest --execute while the connector is stopped.
go deeper
Know you stop the connector first and that there's a REST endpoint to view/clear offsets.
Distinguish source (offset.storage.topic / source maps) from sink (consumer group) reset paths.
Use GET/DELETE/PATCH correctly with the right body shapes per connector type, plus the legacy fallbacks.
Reason about alterOffsets validation, EOS per-connector offsets topics, auto.offset.reset interaction, and safe operational runbooks.
## Step 0: stop the connector — always Never reset offsets on a running connector. A live task can commit offsets concurrently with your reset, racing your change and producing inconsistent state. Use **`PUT /connectors/{name}/stop`** (Kafka 3.5+ added a true STOPPED state that also releases task resources). On older clusters, pause isn't enough for source-topic edits — fully stop/delete or take the cluster action the version supports. ## Modern path: the Connect offsets REST API (KIP-875, Kafka 3.6+) A single, connector-agnostic API handles both source and sink: - **`GET /connectors/{name}/offsets`** — read current offsets (source partition/offset maps, or sink Kafka topic/partition/offset entries). - **`DELETE /connectors/{name}/offsets`** — wipe all committed offsets. A **source** then restarts from its configured initial position (reprocesses from the beginning). A **sink** rewinds and re-consumes per its `auto.offset.reset` (earliest/latest). Connector must be STOPPED. - **`PATCH /connectors/{name}/offsets`** — set **specific** offsets. Body is `{"offsets":[{"partition":{...},"offset":{...}}, ...]}`. - **Source** entry: `partition` = the connector's `sourcePartition` map, `offset` = the `sourceOffset` map (or `null` to clear that one). - **Sink** entry: `partition` = `{"kafka_topic":"t","kafka_partition":0}`, `offset` = `{"kafka_offset":12345}` — because sink offsets are just consumer-group offsets. After resetting, **`PUT /connectors/{name}/resume`** to start consuming/producing again. ## Legacy / pre-3.6 path - **Source connectors:** offsets live in the compacted `offset.storage.topic`, keyed by the serialized source partition. To reset, you produce a record to that topic with the same key and either a new value (set a position) or a null value/tombstone (clear it). This is fiddly and error-prone — exact key serialization must match. Do it only with the connector stopped/deleted. - **Sink connectors:** because the offsets are normal consumer-group offsets, use the standard tooling while the connector is stopped: `kafka-consumer-groups.sh --bootstrap-server ... --group connect-<connector-name> --reset-offsets --to-earliest --all-topics --execute` (or `--to-offset`, `--to-datetime`, `--shift-by`). The group must be inactive, which is why the connector must be stopped. ## Edge cases & gotchas - **Exactly-once source:** each connector has its own offsets topic; use the REST API rather than hand-editing. - **Sink `auto.offset.reset`:** after a DELETE, where a sink restarts (earliest vs latest) is governed by its consumer `auto.offset.reset`, not by Connect itself. - **Partial reset:** PATCH lets you rewind one topic-partition or one source partition without nuking everything — safer than DELETE for surgical reprocessing. - **Validation:** the framework asks the connector to validate offset changes (via `alterOffsets`), so a connector can reject nonsensical positions.
- Why must the connector be stopped before resetting offsets?A running task can commit offsets concurrently, racing your reset and corrupting state. Stopping (or for legacy sinks, leaving the consumer group inactive) ensures your change is the only writer.
- After DELETE-ing a sink connector's offsets, what determines where it restarts?Its consumer auto.offset.reset (earliest or latest). Connect doesn't pick for you; the cleared group falls back to the consumer's reset policy.
- How do you rewind only one topic-partition of a sink without nuking everything?PATCH /connectors/{name}/offsets with a single entry {partition:{kafka_topic,kafka_partition}, offset:{kafka_offset}} — a surgical set instead of a full DELETE.
saying these in an interview costs you the question
- Resetting offsets while the connector is still running.
- Using kafka-consumer-groups --reset-offsets on a SOURCE connector (sources aren't consumer groups).
- Assuming DELETE on a sink always restarts from earliest (it follows auto.offset.reset).
- Hand-editing offset.storage.topic on a 3.6+ cluster instead of using the REST offsets API.