skip to content

Compare SourceRecord and SinkRecord. What fields does each carry, and why does SourceRecord have sourcePartition/sourceOffset while SinkRecord has kafkaPartition/kafkaOffset?

level: seniorimportance: should knowfreq 40%

answer

  1. Both extend ConnectRecord (key/value/schema/headers)
  2. SourceRecord: sourcePartition + sourceOffset (external coords)
  3. SinkRecord: kafkaOffset (Kafka coords)
  4. Source offset -> resume external system
  5. Sink offset -> resume Kafka via consumer group

basics

~20 s

SourceRecord describes data going INTO Kafka: target topic/partition, key, value, schemas, plus sourcePartition/sourceOffset that point back into the external system for resume. SinkRecord describes data coming OUT of Kafka: it carries the actual Kafka topic, kafkaPartition, and kafkaOffset of where it was consumed.

solid answer

~40 s

Both extend ConnectRecord and share topic, kafkaPartition, key, keySchema, value, valueSchema, timestamp, and headers. The difference is the offset semantics that match their direction. A SourceRecord (produced by SourceTask.poll()) additionally carries sourcePartition and sourceOffset — opaque maps that identify a position in the EXTERNAL system (e.g. {"table":"orders"} / {"id":42}). Connect persists these to the connect-offsets topic so the source can resume after restart; the record's topic/partition here are the DESTINATION in Kafka. A SinkRecord (delivered to SinkTask.put()) instead carries kafkaOffset (and originalTopic/originalKafkaPartition) describing exactly where in Kafka the record was read from, because the sink needs that to track consumer offsets and to write partition-aware output. So source offsets look outward into the source system; sink offsets look inward at Kafka's own coordinates.

go deeper

for a junior

Know SourceRecord goes into Kafka and SinkRecord comes out, each with key/value.

for a middle

Identify sourceOffset vs kafkaOffset and which direction each belongs to.

for a senior

Explain why offset coordinates point at the system to resume from, and the role of the offsets topic.

for a principal

Reason about Schema+Struct data model, converter/SMT reuse across directions, and partition-aware sink output keyed by kafkaOffset.

## Shared base: ConnectRecord Both `SourceRecord` and `SinkRecord` extend `ConnectRecord`, which holds the connector-neutral payload: - `topic()` and `kafkaPartition()` - `key()` + `keySchema()` - `value()` + `valueSchema()` (Connect's runtime data model is `Schema` + `Struct`, decoupled from wire format; **converters** translate to/from JSON/Avro/Protobuf) - `timestamp()` - `headers()` Converters and SMTs operate on this common shape, which is why a transform can be reused across source and sink. ## SourceRecord — looking OUTWARD Produced by `SourceTask.poll()`. It represents a record the connector wants to put INTO Kafka, so: - `topic()`/`kafkaPartition()` are the **destination** in Kafka (partition may be null -> the producer's partitioner chooses). - It adds two extra maps: - **sourcePartition** — identifies a logical stream in the external system, e.g. `{"filename":"events.log"}` or `{"server":"db1","table":"orders"}`. - **sourceOffset** — the position within that stream, e.g. `{"position":2048}` or `{"lsn":"0/16B3748"}`. These are **opaque to Connect** — it just stores them in the internal **connect-offsets** topic after the record is acknowledged. On restart, `offsetStorageReader().offset(sourcePartition)` returns the last committed sourceOffset so the task resumes precisely. This is how a CDC connector knows where in the database log to continue. ## SinkRecord — looking INWARD Delivered to `SinkTask.put()`. It represents a record that was CONSUMED from Kafka, so it carries Kafka's own coordinates of where it came from: - `topic()`, `kafkaPartition()`, and crucially **`kafkaOffset()`** — the exact offset in Kafka. - `timestampType()` and the original topic/partition (relevant after SMTs may rewrite topic()). The sink needs `kafkaOffset()` to: (a) report durably-written positions back in `preCommit()` for consumer-offset commit, and (b) do partition-aware output, e.g. an S3 sink writing `topic/partition=3/offset_range.json`. There is **no sourceOffset** on a SinkRecord — the "source" of a sink record *is* Kafka, and Kafka offsets are the system of record. ## Why the asymmetry exists The offset a connector cares about always points at the system it must *resume reading from*: - A source resumes reading the **external system**, so it needs external coordinates -> sourcePartition/sourceOffset. - A sink resumes reading **Kafka**, so it needs Kafka coordinates -> kafkaPartition/kafkaOffset (managed via the consumer group). ## Edge cases / gotchas - A SourceRecord's `kafkaPartition()` being null is normal; the partition is chosen at produce time. - After an SMT like RegexRouter changes a SinkRecord's `topic()`, the original topic is still available, which matters for offset tracking keyed by the real Kafka partition. - Schemas can be null (schemaless mode) — converters then pass through `Map`/primitive values without a Connect Schema.

  • Why is there no sourceOffset on a SinkRecord?
    A sink record's source IS Kafka, so its resume position is the Kafka offset (kafkaOffset). External coordinates only matter for source connectors that resume reading an external system.
  • What runtime type carries the record value inside Connect, independent of JSON/Avro?
    Connect's internal data model: a Schema plus a Struct (or schemaless Map/primitive). Converters serialize this to/from the wire format, so SMTs work on the same shape regardless of converter.

saying these in an interview costs you the question

  • Saying SinkRecord has a sourceOffset.
  • Claiming SourceRecord's topic is the source-system table (it's the Kafka destination topic).
  • Thinking the value is always JSON/Avro inside Connect (it's Schema+Struct; converters handle the wire format).
  • Believing sourcePartition/sourceOffset are interpreted by Connect (they're opaque to it).

context