What is the OffsetStorageReader, and how should a source task use it on startup?
answer
- context.offsetStorageReader() in start()
- offset(partition) / offsets(collection)
- partition map must match poll() exactly
- null => start from initial position
- read-only, reflects last committed flush
basics
~20 sOffsetStorageReader 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.
solid answer
~40 sOffsetStorageReader is the framework-provided read side of source-offset persistence. In SourceTask.start(), the task obtains it via context.offsetStorageReader() and calls offset(sourcePartition) for a single partition or offsets(Collection) for many. It returns the last committed sourceOffset map for each known source partition (or null/empty if none was ever committed). The task uses that to decide where to resume reading the external system — e.g. seek a file to the stored byte position, or set a DB cursor to the stored timestamp/id. It reads from the same store the framework writes to (offset.storage.topic in distributed mode). A correct connector treats a missing offset as 'start from the beginning/configured initial position' and never assumes a value exists.
go deeper
Know it's how a source task resumes from its last saved position on startup.
Explain obtaining it from SourceTaskContext, offset()/offsets(), and the null-means-start-fresh contract.
Diagnose the exact-partition-key-match bug and reason about committed-only visibility.
Discuss the read/write split (reader vs commit cycle), standalone vs distributed backing stores, and resume-correctness guarantees.
## What it is `OffsetStorageReader` is an interface the Connect framework hands to every source task. It is the **read** counterpart to the framework's internal OffsetStorageWriter (which persists offsets during commit cycles). It abstracts away *where* offsets live — `offset.storage.topic` in distributed mode, the offset file in standalone mode — so the connector code is identical in both. ## How a task gets and uses it A `SourceTask` receives a `SourceTaskContext`. In `start(Map<String,String> props)` the task typically does: 1. `OffsetStorageReader reader = context.offsetStorageReader();` 2. Build the `sourcePartition` map(s) it cares about — these MUST be byte-for-byte the same maps the task emits in its SourceRecords, because lookup is by exact key. 3. Call `reader.offset(partition)` (single) or `reader.offsets(partitions)` (batch) to fetch the last committed `sourceOffset` map. 4. If the returned map is non-null, resume from it (seek file position, set DB cursor). If null, fall back to a configured initial position (`snapshot`, `earliest`, a start timestamp, etc.). ## Why the partition key must match exactly Offsets are stored as a key→value where the key is the serialized source partition map. If `start()` builds a partition map that differs even slightly (different key name, different value type, extra field) from what `poll()` emitted, the lookup misses and the task silently restarts from scratch — re-reading everything. This is a classic connector bug. ## Consistency caveats - The reader reflects only **committed** offsets (last flush), so just-produced-but-not-yet-flushed work is not visible — consistent with at-least-once resume semantics. - In distributed mode the reader reads from the compacted offset topic; right after a rebalance the framework ensures the latest committed values are available before `start()`. - It is read-only: a task cannot write offsets through it; persistence happens via the commit cycle on the records it returns. ## Typical mistake Assuming `offset()` never returns null and dereferencing it — a brand-new connector or a new source partition has no stored offset, so the task must handle the null/empty case as 'begin from initial position'.
- A connector restarts and re-reads everything from the beginning even though it ran before. What's a likely cause?The sourcePartition map built in start() doesn't byte-match the one emitted in poll(), so OffsetStorageReader.offset() misses and returns null, making the task start from scratch.
- Can a task use OffsetStorageReader to write or override offsets?No — it's read-only. Offsets are persisted only via the framework's commit cycle on the SourceRecords the task returns from poll().
saying these in an interview costs you the question
- Saying the reader returns uncommitted/in-flight offsets (it reflects the last committed flush only).
- Assuming offset() is never null for a fresh partition.
- Thinking the reader can write offsets back.