Learn Labs
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:

MethodWhenCan it modify?
ProducerRecord<K,V> onSend(ProducerRecord<K,V> record)Before the record is sent to Kafka — indeed before it is even serializedYes. "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 acknowledgmentNo. "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:

  • ProducerInterceptor is a Configurable interface — override configure to 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 (here counting.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.)


On this page