栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 云计算 > 云平台

kafka批量消费

云平台 更新时间: 发布时间: IT归档 最新发布 模块sitemap 名妆网 法律咨询 聚返吧 英语巴士网 伯小乐 网商动力

kafka批量消费

当消费端是批量接收消息,配置中的自动提交需要关闭,同时要把手动提交打开

kafka:
    ###########【Kafka集群】###########
    bootstrap-servers: 192.168.188.128:9092
    producer:
      retries: 0 # 重试次数
      acks: 1 # 应答级别:多少个分区副本备份完成时向生产者发送ack确认(可选0、1、all/-1)
      batch-size: 16384 # 批量大小
      buffer-memory: 33554432 # 生产端缓冲区大小
      # Kafka提供的序列化和反序列化类
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      #额外的,没有直接有properties对应的参数,将存放到下面这个Map对象中,一并初始化
#      properties:
#        #自定义分区器
#        partitioner.class: com.example.demo.kafka.CustomizePartitioner

    consumer:
      group-id: javagroup
      enable-auto-commit: false
      auto-commit-interval: 1000
      # earliest:当各分区下有已提交的offset时,从提交的offset开始消费;无提交的offset时,从头开始消费
      # latest:当各分区下有已提交的offset时,从提交的offset开始消费;无提交的offset时,消费新产生的该分区下的数据
      # none:topic各分区都存在已提交的offset时,从offset后开始消费;只要有一个分区不存在已提交的offset,则抛出异常
      auto-offset-reset: latest
      key-deserializer:  org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer:  org.apache.kafka.common.serialization.StringDeserializer
      max-poll-records: 50
      # 批量消费每次最多消费多少条消息
    listener:
      type: batch
      ack-mode: manual_immediate

注意:批量消费和批量发送是无关的,可以一条条发送批量消费也可以批量发送批量消费
当批量发送时:
Producer:

代码如下:

@Data
public class KafkaMessage {
    public Integer index;
    public String id;
    public String value;
}

@GetMapping("BatchSend")
    public void sendProducerRecord(){
        for(int i=0; i<100; i++) {
            KafkaMessage kafkaMessage = new KafkaMessage();
            kafkaMessage.setIndex(i);
            kafkaMessage.setId(UUID.randomUUID().toString());
            kafkaMessage.setValue("producerRecord " + i);
            ProducerRecord producerRecord =
                    new ProducerRecord<>("topic1", "key1", JSONObject.toJSONString(kafkaMessage));
            kafkaTemplate.send(producerRecord);
        }
    }

Consumer:

代码:

@KafkaListener(id = "consumer1",groupId = "javagroup", topics = "topic1")
    public void onMessage3(List> records) {
        System.out.println(">>>批量消费一次,records.size()="+records.size());
        for (ConsumerRecord record : records) {
            System.out.println("****"+record.value());
        }
    }

ApiPost发送请求:


当Producer发送一条消息,Consumer批量消费时:
Producer:

    @GetMapping("Producer")
    public void sendMessage(String normalMessage){
        kafkaTemplate.send("topic1",normalMessage);
    }

Consumer不变
ApiPost发送请求:


注意:如果把配置文件里批量消费关掉

但是Consumer还是批量接收时,会报错:

因此批量消费时务必修改配置文件中的批量监听以及手动提交
当配置文件中批量监控以及手动提交开启

而Consumer为简单消费时:

无论批量发送还是简单发送都会报错;

总结:
配置关闭,批量消费报错
配置打开,简单消费报错
因此配置和消费要保持一致

当Producer批量发送,配置文件批量监听手动提交都关闭,而Consumer为单体消费时:

@KafkaListener(topics = {"topic1"})
    public void onMessage1(ConsumerRecordrecord){
        System.out.println("简单消费: "+record.topic()+"_"+record.partition()+"_"+"_"+record.value());
    }


二.回调:
kafkaTemplate提供了一个回调方法addCallback,我们可以在回调方法中监控消息是否发送成功 或 失败时做补偿处理,有两种写法,

//    带回掉的Producer  .addCallback(SuccessCallback<>,FailureCallback<>)
    @GetMapping("CallBackOne")
    public void sendMessage2(String callBackMessage){
        kafkaTemplate.send("topic1",callBackMessage).addCallback(successCallBack->{
            String topic = successCallBack.getRecordMetadata().topic();
            int partition = successCallBack.getRecordMetadata().partition();
            long offset = successCallBack.getRecordMetadata().offset();
            System.out.println("发送消息成功: "+topic+"-"+partition+"-"+offset);
        },failureCallBack->{
            System.out.println("发送消息失败:"+failureCallBack.getMessage());
        });
    }


第二种写法:

//带回调的Producer addCallback(new ListenableFutureCallback>() {})
    @GetMapping("CallBackTwo")
    public void sendMessage3(String callBackMessage){
        kafkaTemplate.send("topic2",callBackMessage).addCallback(new ListenableFutureCallback>() {
            @Override
            public void onFailure(Throwable ex) {
                System.out.println("发送消息失败: "+ex.getMessage());
            }

            @Override
            public void onSuccess(SendResult result) {
                System.out.println("发送消息成功: "+"-"+result.getRecordMetadata().offset());
            }
        });
    }


补充:全局回调
注意:这样设置以后,该单例的kafkaTemplate就有了一个回调,注意是全局回调,只要一个实例中设置,全局生效

@Component
@Slf4j
public class KafkaSendResultHandler implements ProducerListener {

    @Override
    public void onSuccess(ProducerRecord producerRecord, RecordMetadata recordMetadata) {
        System.out.println("Message send success : " + producerRecord.toString());

    }

    @Override
    public void onError(ProducerRecord producerRecord, RecordMetadata recordMetadata, Exception exception) {
        System.out.println("Message send error : " + producerRecord.toString());
    }
}
    @Autowired(required = false)
    private KafkaSendResultHandler producerListener;
    //全局回调
    @GetMapping("GlobalCallBack")
    public void testProducerListen() throws InterruptedException {
        kafkaTemplate.setProducerListener(producerListener);
        kafkaTemplate.send("topic1", "test producer listen");
        Thread.sleep(1000);
    }


此时如果调用另一个接口:


此接口虽然未设置回调,但由于是全局回调,所以回调依然生效

三:自定义分区器

我们知道,kafka中每个topic被划分为多个分区,那么生产者将消息发送到topic时,具体追加到哪个分区呢?这就是所谓的分区策略,Kafka 为我们提供了默认的分区策略,同时它也支持自定义分区策略。其路由机制为:

① 若发送消息时指定了分区(即自定义分区策略),则直接将消息append到指定分区;

② 若发送消息时未指定 patition,但指定了 key(kafka允许为每条消息设置一个key),则对key值进行hash计算,根据计算结果路由到指定分区,这种情况下可以保证同一个 Key 的所有消息都进入到相同的分区;

③ patition 和 key 都未指定,则使用kafka默认的分区策略,轮询选出一个 patition;

※ 我们来自定义一个分区策略,将消息发送到我们指定的partition,首先新建一个分区器类实现Partitioner接口,重写方法,其中partition方法的返回值就表示将消息发送到几号分区,

public class CustomizePartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        // 得到 topic 的 partitions 信息
        List partitions = cluster.partitionsForTopic(topic);
        int numPartitions = partitions.size();
        // 模拟某客服
        if(key.toString().equals("10000") || key.toString().equals("11111")) {
            // 放到最后一个分区中
            return numPartitions - 1;
        }
        String phoneNum = key.toString();
        System.out.println("phoneNum"+phoneNum);
        return phoneNum.substring(0, 3).hashCode() % (numPartitions);
    }

    @Override
    public void close() {

    }

    @Override
    public void configure(Map map) {

    }
}

由于只划了一个分区,所以消息实际上全都分配到同一个分区

转载请注明:文章转载自 www.mshxw.com
本文地址:https://www.mshxw.com/it/895722.html
我们一直用心在做
关于我们 文章归档 网站地图 联系我们

版权所有 (c)2021-2022 MSHXW.COM

ICP备案号:晋ICP备2021003244-6号