A Java interceptor sees each record in onSend() and each answer in onAcknowledgement() without touching application code. listings/l0657_plugins.sh compiles it and the partitioner, puts the jar on the console producer's CLASSPATH, and sets partitioner.class and interceptor.classes.
import java.util.Map;
import java.util.concurrent.atomic.AtomicLong;
import org.apache.kafka.clients.producer.*;
public class AuditInterceptor implements ProducerInterceptor<byte[], byte[]> {
private final AtomicLong acked = new AtomicLong(), failed = new AtomicLong();
@Override public ProducerRecord<byte[], byte[]> onSend(ProducerRecord<byte[], byte[]> r) {
r.headers().add("produced-by", "booknest-shop".getBytes()); // before partitioning
return r;
}
@Override public void onAcknowledgement(RecordMetadata m, Exception e) {
(e == null ? acked : failed).incrementAndGet(); // I/O thread: be quick
}
@Override public void close() {
System.out.printf("audit: %d acknowledged, %d failed%n", acked.get(), failed.get());
}
@Override public void configure(Map<String, ?> configs) {}
}Output
audit: 3 acknowledged, 0 failed
Partition:0 produced-by:booknest-shop 9 {}
Partition:1 produced-by:booknest-shop 7 {}
Partition:2 produced-by:booknest-shop 8 {}Java producers also publish JMX metrics (Monitoring and Kafka UIs), notably record-send-rate, record-error-rate, request-latency-avg and batch-size-avg; librdkafka emits JSON statistics, which producer_lab.py reads.