What is a ProducerInterceptor, and what do onSend and onAcknowledgement do, including thread and ordering semantics?
answer
- interceptor.classes = ordered chain
- onSend: caller thread, pre-serialization, returns record
- onAcknowledgement: Sender/IO thread, must not block
- exactly one of metadata/exception non-null
- exceptions caught + logged, send continues
basics
~10 sA ProducerInterceptor lets you hook into the producer pipeline. onSend runs before serialization and can mutate/inspect the record; onAcknowledgement runs when the broker acks or the send fails. You configure a chain via interceptor.classes.
solid answer
~40 sProducerInterceptor<K,V> (configured via the interceptor.classes list) intercepts records on their way out. onSend(ProducerRecord) is called on the application thread inside send(), before key/value serialization and partition assignment; it can return a modified record (e.g. add headers, tag for tracing) — but mutating partition/key changes routing, so it's used carefully. onAcknowledgement(RecordMetadata, Exception) fires when the record is acknowledged by the broker or fails; exactly one of metadata/exception is non-null. It runs on the producer's I/O (Sender) thread, so it must be fast and non-blocking or it stalls all sends. Interceptors form a chain in config order: onSend runs first-to-last, and each interceptor sees the previous one's output. Exceptions thrown by an interceptor are caught and logged — they don't break the send chain. Common uses: distributed tracing (inject trace headers), metrics/auditing, and Confluent's monitoring interceptor.
go deeper
Know a ProducerInterceptor hooks the send pipeline with onSend and onAcknowledgement, set via interceptor.classes.
Explain that onSend can modify the record pre-serialization and onAcknowledgement fires on success or failure with one of metadata/exception.
Discuss thread semantics (caller vs Sender thread), the non-blocking requirement, chain ordering, and swallowed exceptions.
Design tracing/audit infra via interceptors org-wide while guarding the Sender thread; know when a serializer is the right hook instead.
## Purpose A `ProducerInterceptor<K,V>` is a pluggable hook in the producer's send pipeline, used for **cross-cutting concerns** — tracing, metrics, auditing, header injection — without touching business code. Configure a chain: ``` interceptor.classes=com.acme.TracingInterceptor,com.acme.AuditInterceptor ``` ## The interface ``` public interface ProducerInterceptor<K,V> extends Configurable, AutoCloseable { ProducerRecord<K,V> onSend(ProducerRecord<K,V> record); void onAcknowledgement(RecordMetadata metadata, Exception exception); void close(); } ``` ### onSend - Called **on the application/caller thread** from inside `KafkaProducer.send()`, **before serialization and before the partitioner runs**. - Receives the record and **returns a (possibly new) record**. Returning a modified record lets you add headers or rewrite fields. The returned record from one interceptor is the input to the next. - It sees the record **pre-serialization** (objects, not bytes), so it can read the typed key/value. Changing the key or target partition **alters routing** — generally you only add headers. - If `onSend` throws, the exception is **caught and logged**; the original (pre-interceptor) record continues — interceptors must not break sends. ### onAcknowledgement - Called when the record is **acknowledged by the broker** or when the send **fails** (e.g. serialization error, buffer exhaustion, broker error). - Exactly one of `(metadata, exception)` is non-null: success → metadata with topic/partition/offset; failure → exception, and metadata has partition info but offset -1. - Runs on the producer **I/O (Sender) thread** — the same thread that batches and transmits. **It must be fast and non-blocking**: any blocking call (network, lock, heavy compute) here stalls *all* in-flight sends for the producer. It also runs *before* the user's `Callback`/`Future` completes, so heavy work here delays delivery callbacks. ## Chain ordering - `interceptor.classes` defines order. `onSend` executes **first→last**, threading each output into the next. - `onAcknowledgement` is invoked on every interceptor as well; an exception in one is caught/logged and does not stop the others. ## Lifecycle - `configure(Map)` (from `Configurable`) runs at producer startup; `close()` at producer close. - Interceptors are **shared across all sends** of the producer instance, so they must be thread-safe if they hold state (onSend on caller threads, onAcknowledgement on the Sender thread can run concurrently). ## Common real uses - **Distributed tracing**: OpenTelemetry's `TracingProducerInterceptor` injects W3C traceparent headers in onSend. - **Confluent Monitoring Interceptor**: records produce metrics for Control Center. - **Auditing/metrics**: count records, measure end-to-end ack latency in onAcknowledgement. ## Edge cases / gotchas - Blocking in onAcknowledgement is the classic production incident — it serializes onto the single Sender thread and tanks throughput. - onSend running before serialization means an interceptor can't see or fix the serialized bytes — for that you'd write a serializer, not an interceptor. - Because onSend exceptions are swallowed, a buggy interceptor can silently no-op; test it explicitly. - Interceptors do not run for records dropped before send (they aren't a validation layer you can rely on to reject).
- Why must onAcknowledgement be non-blocking?It runs on the producer's single I/O (Sender) thread that batches and transmits all records. Blocking there stalls every in-flight send and delays delivery callbacks, crushing throughput.
- Does onSend see the serialized bytes?No. onSend runs before serialization and the partitioner, so it sees the typed key/value objects. To touch the bytes you'd implement a Serializer instead.
- What happens if an interceptor's onSend throws?The exception is caught and logged; the send proceeds with the original record. Interceptors are not allowed to break the send path.
saying these in an interview costs you the question
- Saying onAcknowledgement runs on the caller/application thread (it runs on the Sender I/O thread)
- Claiming onSend runs after serialization (it runs before serialization and partitioning)
- Thinking an interceptor exception aborts the send (it's caught and logged)
- Treating interceptors as a reliable validation/rejection mechanism