skip to content

In what order do interceptors, serializers, and the partitioner execute in the producer send path, and what does each operate on?

level: middleimportance: should knowfreq 30%

answer

  1. onSend(objects) → serialize(bytes) → partition(keyBytes) → accumulate → send → onAck
  2. serialization = object/byte boundary
  3. default partitioner hashes serialized key (murmur2)
  4. null key → sticky partitioning
  5. onAck on Sender thread

basics

~10 s

Order inside send(): interceptor.onSend (sees objects) -> key/value serializers (produce bytes) -> partitioner (uses serialized key) -> accumulator/batching -> network. onAcknowledgement fires later on ack/failure.

solid answer

~40 s

When you call producer.send(record), on the calling thread the producer first runs each ProducerInterceptor.onSend in chain order — these see the typed key/value objects and can add headers. Then the configured key.serializer and value.serializer convert the (possibly modified) key and value to byte[]. Next the partitioner picks a partition; the default partitioner hashes the serialized key bytes (or uses sticky partitioning for null keys), so it operates on bytes, not the original object. The record then enters the RecordAccumulator where it's batched by partition, and the Sender thread transmits it. Much later, when the broker acks (or the send fails), onAcknowledgement runs on the Sender thread. The key insight: interceptors see objects pre-serialization, the partitioner sees serialized key bytes, and serialization is therefore a hard boundary between them.

go deeper

for a junior

Know the rough order: interceptor, then serialize, then choose partition, then send.

for a middle

State that onSend sees objects, serializers make bytes, the default partitioner hashes the serialized key, all on the caller thread.

for a senior

Explain the object/byte boundary, sticky partitioning for null keys, and that onAcknowledgement runs later on the Sender thread.

for a principal

Reason about throughput: serialization cost on the caller thread, where to inject routing logic, and interceptor placement in the pipeline.

## The send pipeline, step by step Calling `producer.send(record, callback)` does the following **on the application/caller thread**: 1. **Interceptor chain — onSend**: each `ProducerInterceptor.onSend` runs in `interceptor.classes` order, receiving the typed `ProducerRecord<K,V>` (objects, not bytes). Output threads into the next interceptor. Used for headers, tracing, tagging. 2. **Serialization**: `key.serializer.serialize(topic, headers, key)` then `value.serializer.serialize(topic, headers, value)` convert to `byte[]`. A failure here throws `SerializationException` synchronously from `send()` (not retried). 3. **Partitioning**: the `Partitioner` chooses a partition. The default (`DefaultPartitioner`/built-in logic) **hashes the serialized key bytes** (murmur2) when a key is present; with a null key it uses **sticky partitioning** (batches to one partition until it fills, then rotates) for better batching. So the partitioner sees **bytes**, and a custom partitioner reading the original object isn't possible at this stage — it gets the serialized key. 4. **Accumulation**: the record is appended to a per-partition batch in the `RecordAccumulator` (in-memory buffer, sized by `buffer.memory`, batched by `batch.size`/`linger.ms`). 5. **Transmission**: the background **Sender** thread drains ready batches and sends them to broker leaders, honoring `acks`, `max.in.flight.requests.per.connection`, retries, and idempotence. 6. **Acknowledgement**: when the broker responds (or the send ultimately fails), the **Sender thread** invokes the user `Callback`, completes the `Future<RecordMetadata>`, and calls each interceptor's `onAcknowledgement`. ## What each stage operates on | Stage | Thread | Sees | |---|---|---| | onSend | caller | typed objects (K,V) | | serializers | caller | objects → bytes | | partitioner | caller | serialized key bytes | | accumulator | caller | byte batches | | Sender / onAcknowledgement | Sender I/O | RecordMetadata/Exception | ## Why the order matters - An interceptor that wants to **influence partitioning by a business field** must set the key (or partition) in onSend, because by the time the partitioner runs only bytes remain. - Serialization is the **hard boundary**: anything object-level (interceptors, your callback's view of the record) happens before it; anything byte-level (partitioner hash, compression, batching) after. - Serialization runs on the **caller** thread, so it's a cost you pay synchronously; heavy serializers slow your producing loop. ## Edge cases - Headers added in onSend are visible to the serializer's `serialize(topic, headers, data)` overload and are sent with the record. - If onSend changes the value to a type incompatible with value.serializer, you get a ClassCastException at step 2. - A custom partitioner can still inspect the original key/value object via the `partition(topic, key, keyBytes, value, valueBytes, cluster)` signature — it receives both, but the default partitioner uses keyBytes.

  • Can an interceptor influence which partition a record lands on?
    Yes, but only by setting the key or explicit partition in onSend, because the partitioner runs after serialization and (for the default) hashes the serialized key bytes — it can't see the original object's business fields otherwise.
  • On which thread does serialization happen and why does it matter?
    On the application/caller thread inside send(). Heavy serialization is paid synchronously by your producing loop, so an expensive serializer directly limits per-thread throughput.

saying these in an interview costs you the question

  • Putting serialization before interceptors (onSend runs first, on objects)
  • Saying the partitioner hashes the original key object (the default hashes serialized key bytes)
  • Claiming serialization runs on the Sender thread (it runs on the caller thread)
  • Forgetting that null keys trigger sticky partitioning, not round-robin per-record

context