skip to content

Walk through the SourceTask.poll() loop and the SinkTask.put()/flush() loop. What is each responsible for, and how do offsets get committed?

level: seniorimportance: must knowfreq 60%

answer

  1. poll() returns SourceRecords; block, don't spin
  2. sourcePartition + sourceOffset embedded per record
  3. Source offsets -> connect-offsets topic, not consumer offsets
  4. put() writes/buffers; preCommit() flushes + returns safe offsets
  5. Write downstream THEN commit Kafka offset = at-least-once

basics

~20 s

SourceTask.poll() returns new SourceRecords for Connect to produce into Kafka; Connect tracks the source offsets you embed. SinkTask.put() receives batches of SinkRecords to write downstream; flush()/preCommit() is where the sink ensures records are durable before Kafka offsets are committed.

solid answer

~50 s

On the source side, Connect calls poll() in a loop on a dedicated thread; it returns a List<SourceRecord>, each carrying a sourcePartition and sourceOffset (e.g. {file: a.log} / {position: 1024}). Connect produces these to Kafka and, after they're acknowledged, persists the source offsets to the internal connect-offsets topic, so a restart resumes from the right place. poll() should block until data is available rather than busy-spin. On the sink side, Connect's consumer fetches records and hands them to put(Collection<SinkRecord>) in batches; the task buffers or writes them downstream. Connect periodically calls flush() (or preCommit()) — the task must ensure all records up to the given offsets are durably written to the external system, then Connect commits the corresponding Kafka consumer offsets. preCommit() lets the task return the exact offsets it has safely persisted, enabling correct at-least-once (or, with idempotent sinks, effectively-once) semantics.

go deeper

for a junior

Know poll() produces records into Kafka and put() consumes records out.

for a middle

Describe sourcePartition/sourceOffset and that flush/preCommit precedes offset commit.

for a senior

Explain the write-then-commit ordering, preCommit() return semantics, and at-least-once vs exactly-once.

for a principal

Reason about exactly-once source support (KIP-618), idempotent/transactional sinks, and offset-store-with-data patterns for effectively-once.

## Source side: the poll() loop A `SourceTask` runs on its own thread inside a worker. The framework loops: ``` while (running) { List<SourceRecord> records = task.poll(); // your code for (record : records) producer.send(record); // after acks, commit source offsets } ``` ### What poll() must do - Return a `List<SourceRecord>` of *new* data fetched from the external system since last time. - **Block** when there's nothing new (e.g. wait on a queue / sleep briefly), returning null or an empty list rather than spinning hot. - Embed two maps in each record: **sourcePartition** (identifies a logical stream, e.g. `{"filename":"a.log"}`) and **sourceOffset** (position within it, e.g. `{"position":1024}`). ### Source offset commit Connect does NOT use Kafka consumer offsets for sources (there's no consumer). Instead it produces the records to the destination topic, and once the producer acknowledges them, it writes the (sourcePartition -> sourceOffset) pairs to the internal **connect-offsets** topic (`offset.storage.topic`). On restart, `SourceTaskContext.offsetStorageReader()` lets the task read back the last committed offset and resume. `offset.flush.interval.ms` controls how often offsets are flushed. `commitRecord()` is an optional callback after each record is acked. ## Sink side: the put()/flush() loop A `SinkTask` is backed by a Kafka **consumer** managed by Connect. The framework loops: ``` while (running) { ConsumerRecords r = consumer.poll(); task.put(toSinkRecords(r)); // your code: write/buffer downstream if (timeToCommit) { offsets = task.preCommit(currentOffsets); // your code: flush + return safe offsets consumer.commitSync(offsets); } } ``` ### put(Collection<SinkRecord>) - Receives a batch of `SinkRecord`s. The task writes them to the external system, or buffers them for batched writes. - Must handle retries; throwing `RetriableException` makes Connect retry the same batch. ### flush() / preCommit() - On a periodic boundary (`offset.flush.interval.ms`), Connect calls `flush(currentOffsets)` or, preferably, `preCommit(currentOffsets)`. - The task must **make all buffered records durable downstream** before returning. - `preCommit()` returns the map of `TopicPartition -> OffsetAndMetadata` that the task has *actually* persisted. Connect commits exactly those to Kafka. Returning a lower offset than `currentOffsets` is how a task says "I've only safely written up to here." Returning an empty map disables Connect-managed offset commit (used when the sink stores offsets itself, e.g. for exactly-once into the external store). ## Why the order matters (delivery semantics) The golden rule: **persist downstream, THEN commit the Kafka offset.** If the task committed first and then crashed before writing, data would be lost (less than at-least-once). By writing first and committing the safe offset in preCommit(), Connect guarantees **at-least-once**. To reach **effectively/exactly-once**, the sink must be idempotent or transactional (e.g. keyed upserts, or storing offsets atomically with the data). For sources, exactly-once is available via the dedicated exactly-once source support (KIP-618: `exactly.once.source.support`) which wraps producer transactions around record batches and offset writes. ## Common pitfalls - Busy-spinning in poll() instead of blocking -> CPU burn. - Committing Kafka offsets in put() rather than after durable write -> data loss window. - Ignoring RetriableException vs a fatal exception -> unnecessary task FAILED state. - Not implementing preCommit() so offsets advance past not-yet-flushed data.

  • Why does preCommit() return a map of offsets instead of just void?
    So the task can tell Connect the exact offsets it has durably persisted. Connect commits only those; returning lower offsets prevents committing past un-flushed data, and an empty map disables Connect-managed commits when the sink manages offsets itself.
  • How do source offsets survive a restart given there's no consumer group?
    Connect writes (sourcePartition->sourceOffset) to the internal connect-offsets topic after records are acked. On restart the task reads them back via offsetStorageReader() and resumes.

saying these in an interview costs you the question

  • Saying source connectors use Kafka consumer offsets (they don't — they use the connect-offsets topic).
  • Committing the Kafka offset before the downstream write completes.
  • Claiming put() must return offsets — that's preCommit()/flush().
  • Busy-looping in poll() instead of blocking until data arrives.

context