Compare SourceRecord and SinkRecord. What fields does each carry, and why does SourceRecord have sourcePartition/sourceOffset while SinkRecord has kafkaPartition/kafkaOffset?
answer
- Both extend ConnectRecord (key/value/schema/headers)
- SourceRecord: sourcePartition + sourceOffset (external coords)
- SinkRecord: kafkaOffset (Kafka coords)
- Source offset -> resume external system
- Sink offset -> resume Kafka via consumer group
basics
~20 sSourceRecord 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 sBoth 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
Know SourceRecord goes into Kafka and SinkRecord comes out, each with key/value.
Identify sourceOffset vs kafkaOffset and which direction each belongs to.
Explain why offset coordinates point at the system to resume from, and the role of the offsets topic.
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).