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?
answer
- flush.interval.ms default 60000
- produce -> ack -> flush -> write offsets -> commit()
- offsets committed AFTER data = at-least-once
- flush.timeout.ms default 5000 aborts cycle
- interval bounds duplicate window, not throughput
basics
~20 sA 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.
solid answer
~40 sThe source task's poll() returns SourceRecords, each carrying a sourcePartition/sourceOffset. The framework's WorkerSourceTask sends those records via an internal producer and tracks which are acknowledged. On the offset.flush.interval.ms timer (default 60s), the framework runs a commit cycle: it flushes the producer so all outstanding records are durably written and acked, then writes the highest acknowledged source offset per source partition to offset.storage.topic, then calls the task's commit()/commitRecord() callbacks. The ordering guarantees offsets are never committed ahead of data — so a crash before flush replays from the last committed offset, giving at-least-once delivery (duplicates possible, no data loss). If a flush can't finish within offset.flush.timeout.ms, the commit is aborted and retried next cycle.
go deeper
Know the default (60s) and that offsets are saved periodically after data is in Kafka.
Explain the produce->ack->flush->write-offsets ordering and why it yields at-least-once with a duplicate window sized by the interval.
Discuss offset.flush.timeout.ms aborts, callback semantics (commit/commitRecord acking the source), and tuning trade-offs.
Contrast the timer-based cadence with transaction-boundary commits under exactly-once source and reason about failure-mode guarantees end to end.
## The actors - **SourceTask** — your connector code. Its `poll()` returns a `List<SourceRecord>`; each record carries the destination topic/value plus a `sourcePartition` and `sourceOffset` describing where in the external system it came from. - **WorkerSourceTask** — the framework wrapper around your task. It owns an internal **producer** to Kafka and an **OffsetStorageWriter**. ## The flow, step by step 1. The framework calls `poll()`; you return records. 2. The framework sends each record through its producer to the target Kafka topic. Sends are asynchronous; the framework records the send and registers a callback that fires when the broker **acks** the record. 3. As acks arrive, the framework remembers, per source partition, the highest source offset that is now durably in Kafka. 4. On a timer governed by **`offset.flush.interval.ms`** (default **60000 ms**), a **commit cycle** runs: a. **flush the producer** — block until all outstanding sends are acked (bounded by `offset.flush.timeout.ms`, default 5000 ms). b. **write source offsets** — the OffsetStorageWriter writes the latest acked source offsets to `offset.storage.topic` (a compacted topic) and flushes that write. c. **invoke task callbacks** — `commit()` and per-record `commitRecord()` let the connector ack back to the external system (e.g. advance a DB log position, delete an SQS message). ## Why the ordering is sacred Offsets are committed **after** data is acknowledged, never before. Consequences: - Crash **after** producing but **before** committing offsets -> on restart the task resumes from the last *committed* offset and **re-emits** the un-committed records = **duplicates** = **at-least-once** semantics (the default). - Offsets are never ahead of data, so you can't silently **lose** records. ## Tuning trade-offs - **Smaller `offset.flush.interval.ms`** -> offsets committed more often -> fewer duplicates after a crash, but more overhead (more flushes, more writes to the offset topic). - **Larger interval** -> less overhead, but a crash can replay up to a whole interval's worth of records. - **`offset.flush.timeout.ms`** caps how long the flush may take; if exceeded, the framework logs and **skips** committing this cycle (the records stay un-committed and will be retried/replayed). Persistent timeouts usually mean the producer or offset topic is slow/unavailable. ## Important nuance `offset.flush.interval.ms` does **not** throttle data production — records keep flowing to Kafka continuously. It only governs how often the *progress bookmark* is persisted. With exactly-once source enabled the commit cadence is instead tied to **transaction boundaries** (per batch / per connector-defined boundary), not this timer.
- If a worker crashes 30 seconds into a 60-second flush interval, what happens to the records produced in those 30 seconds?Their offsets were never committed, so on restart the task resumes from the last committed offset and re-produces them — duplicates downstream (at-least-once).
- Does lowering offset.flush.interval.ms slow down throughput?No. Production is continuous; the interval only controls how often progress is persisted. Lowering it shrinks the replay/duplicate window at the cost of more offset-topic writes.
saying these in an interview costs you the question
- Saying offsets are committed before/at the same time as producing (they're committed only after acks).
- Claiming offset.flush.interval.ms throttles record production.
- Thinking Connect source connectors are exactly-once by default (they're at-least-once unless exactly.once.support is enabled).