kafka生产者的数据生产流程中,有三个环节是我们可以自定义的,如下图所示。本文为大家介绍如何自定义kafka生产者拦截器。

通过实现kafka提供的拦截器接口类ProducerInterceptor,可以实现消息拦截器的效果。拦截器接口类及方法功能详见下方的注释。
package org.apache.kafka.clients.producer;
import org.apache.kafka.common.Configurable;
//
public interface ProducerInterceptor<K, V> extends Configurable {
/**
* 该方法封装于KafkaProducer.send()方法中,运行在用户主线程
* Producer确保在消息序列化前调用该方法,可以对消息进行任意操作,但慎重修改消息的topic、key和partition,会影响分区以及日志压缩
*/
public ProducerRecord<K, V> onSend(ProducerRecord<K, V> record);
/**
* 该方法在消息发送结果应答或者发送失败时调用,并且通常都是在callback()触发之前执行,运行在IO线程中
。实现该方法的代码逻辑尽量简单,否则影响消息发送效率
*/
public void onAcknowledgement(RecordMetadata metadata, Exception exception);
/**
* This is called when interceptor is closed
*/
public void close();
}
public class MyProducerInterceptor implements ProducerInterceptor {
AtomicInteger successCnt = new AtomicInteger(); //发送成功消息计数器
AtomicInteger failureCnt = new AtomicInteger(); //发送失败消息计数器
@Override
public ProducerRecord onSend(ProducerRecord record) {
//在这里可以改变消息内容
Object newValue = record.value();
return new ProducerRecord(record.topic(),record.partition(),record.timestamp(),record.key(),newValue);
}
@Override
public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
if (metadata == null) {
failureCnt.getAndIncrement(); //失败计数器加1
}else {
successCnt.getAndIncrement(); //成功计数器加1
}
}
/**
* 生产者的producer.close触发
*/
@Override
public void close() {
double successRate = (double) successCnt.get() / (successCnt.get() + failureCnt.get());
System.out.println("消息发送成功率:" + successRate*100 +"%");
}
@Override
public void configure(Map<String, ?> map) {
}
}
生产者拦截器可做消息发送前以及producer回调前的定制化需求,允许用户指定多个Interceptor按照配置顺序作用于一条消息从而形成一个拦截链
List<String> interceptors = new ArrayList<>();
interceptors.add("com.zimug.producer.TestAInterceptor");
interceptors.add("com.zimug.producer.TestBInterceptor");
props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, interceptors);
也可以单独指定配置一个拦截器
//通过这一行生产者配置参数,指定了我们自定义的生产者拦截器
props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, MyProducerInterceptor.class.getName());
测试上文中消息发送成功率统计的方法,如无意外发生,当下文代码中的producer.close();调用时打印消息发送成功率:100.0%,表示我们的生产者消息拦截器生效了。
//没有回调函数的调用方法
@Test
public void noCallback() {
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 20; i++) {
producer.send(
new ProducerRecord<>("producer_test",Integer.toString(i),"noCallback value:" + i)
);
}
producer.close();
}