• 【kafka专栏】自定义生产者客户端拦截器-实现消息发送成功率统计


    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();
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22

    二、拦截器实现:消息发送成功率统计

    • 定义两个计数器,如下文代码所示。
    • onAcknowledgement方法在消息发送结果应答或者发送失败时调用,所以我们个可以根据回调函数参数metadata判断消息发送是否成功
    • 在OnSend方法里面可以改变消息数据,比如:可以在下文为字符串消息加了一个前缀。
    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) {
    
        }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31

    三、指定生产者拦截器

    生产者拦截器可做消息发送前以及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);
    
    • 1
    • 2
    • 3
    • 4

    也可以单独指定配置一个拦截器

    //通过这一行生产者配置参数,指定了我们自定义的生产者拦截器
    props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, MyProducerInterceptor.class.getName());
    
    • 1
    • 2

    测试上文中消息发送成功率统计的方法,如无意外发生,当下文代码中的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();
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
  • 相关阅读:
    数据删掉了怎么恢复?数据删除后还能恢复吗
    视频智能分析国标GB28181云平台EasyCVR加密机授权异常是什么原因?
    springboot微信点餐系统的设计与实现毕业设计源码221541
    C++ ++ 和 -- 运算符重载
    《C++ Primer》第5章 语句
    Webpack5入门到原理
    dockerfile lnmp 搭建wordpress、docker-compose搭建wordpress
    Promes 基于飞书的机器人告警推送
    第十六章 Spring Cloud Alibaba 基础环境搭建
    tqdm使用
  • 原文地址:https://blog.csdn.net/hanxiaotongtong/article/details/125532848