3. Kafka Producers: Writing Messages to Kafka
3.11 Interceptors
ProducerInterceptor has two key methods:
The problem they solve:
"There are times when you want to modify the behavior of your Kafka client application without modifying its code, perhaps because you want to add identical behavior to all applications in the organization. Or perhaps you don't have access to the original code."
ProducerInterceptor has two key methods:
| Method | When | Can it modify? |
|---|---|---|
ProducerRecord<K,V> onSend(ProducerRecord<K,V> record) | Before the record is sent to Kafka — indeed before it is even serialized | Yes. "you can capture information about the sent record and even modify it. Just be sure to return a valid ProducerRecord — the record this method returns will be serialized and sent to Kafka." |
void onAcknowledgement(RecordMetadata metadata, Exception exception) | If and when Kafka responds with an acknowledgment | No. "does not allow modifying the response from Kafka, but you can capture information about the response." |
Common use cases:
- capturing monitoring and tracing information
- enhancing the message with standard headers, especially for lineage tracking
- redacting sensitive information
Example: a counting interceptor
public class CountingProducerInterceptor implements ProducerInterceptor {
ScheduledExecutorService executorService =
Executors.newSingleThreadScheduledExecutor();
static AtomicLong numSent = new AtomicLong(0);
static AtomicLong numAcked = new AtomicLong(0);
public void configure(Map<String, ?> map) {
Long windowSize = Long.valueOf(
(String) map.get("counting.interceptor.window.size.ms"));
executorService.scheduleAtFixedRate(CountingProducerInterceptor::run,
windowSize, windowSize, TimeUnit.MILLISECONDS);
}
public ProducerRecord onSend(ProducerRecord producerRecord) {
numSent.incrementAndGet();
return producerRecord; // unmodified
}
public void onAcknowledgement(RecordMetadata recordMetadata, Exception e) {
numAcked.incrementAndGet();
}
public void close() {
executorService.shutdownNow(); // clean up! avoid leaks
}
public static void run() {
System.out.println(numSent.getAndSet(0));
System.out.println(numAcked.getAndSet(0));
}
}Notes from the book:
ProducerInterceptoris aConfigurableinterface — overrideconfigureto set up before any other method is called. It receives the entire producer configuration, so you can read any parameter, including your own custom ones (herecounting.interceptor.window.size.ms).close()is "the place to close everything and avoid leaks" — threads, file handles, connections to remote data stores.
Applying an interceptor with zero code changes
Using it with kafka-console-producer (which ships with Kafka):
# 1. add your jar to the classpath
export CLASSPATH=$CLASSPATH:~./target/CountProducerInterceptor-1.0-SNAPSHOT.jar
# 2. create a config file containing:
# interceptor.classes=com.shapira.examples.interceptors.CountProducerInterceptor
# counting.interceptor.window.size.ms=10000
# 3. run normally, including the config
bin/kafka-console-producer.sh --broker-list localhost:9092 \
--topic interceptor-test --producer.config producer.config(Note: --broker-list here is the older flag form; Ch. 2 covers the --bootstrap-server migration.)