Day26 MQ补充Kafka相关

客户端消息流转流程

image.png

  • 生产者发送主流程

    1. 设置Producer核心属性:Producer可选的属性都可以由ProducerConfig类管理
    2. 构建消息:构建Key-Value结构的消息,key和value都可以是任意对象类型,其中key主要是用来进行Partition分区的,业务上更关系value
    3. 发送消息:常用的有单向发送、同步发送、异步发送三种方式
    java
    复制代码
    public class MyProducer { private static final String BOOTSTRAP_SERVERS = "worker1:9092,worker2:9092,worker3:9092"; private static final String TOPIC = "disTopic"; public static void main(String[] args) throws ExecutionException, InterruptedException { //PART1:设置发送者相关属性 Properties props = new Properties(); // 此处配置的是kafka的端口 props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS); // 配置key的序列化类 props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer"); // 配置value的序列化类 props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer"); Producer<String,String> producer = new KafkaProducer<>(props); CountDownLatch latch = new CountDownLatch(5); for(int i = 0; i < 5; i++) { //Part2:构建消息 ProducerRecord<String, String> record = new ProducerRecord<>(TOPIC, Integer.toString(i), "MyProducer" + i); //Part3:发送消息 //单向发送:不关心服务端的应答。 producer.send(record); System.out.println("message "+i+" sended"); //同步发送:获取服务端应答消息前,会阻塞当前线程。 RecordMetadata recordMetadata = producer.send(record).get(); String topic = recordMetadata.topic(); int partition = recordMetadata.partition(); long offset = recordMetadata.offset(); String message = recordMetadata.toString(); System.out.println("message:["+ message+"] sended with topic:"+topic+"; partition:"+partition+ ";offset:"+offset); //异步发送:消息发送后不阻塞,服务端有应答后会触发回调函数 producer.send(record, new Callback() { @Override public void onCompletion(RecordMetadata recordMetadata, Exception e) { if(null != e){ System.out.println("消息发送失败,"+e.getMessage()); e.printStackTrace(); }else{ String topic = recordMetadata.topic(); long offset = recordMetadata.offset(); String message = recordMetadata.toString(); System.out.println("message:["+ message+"] sended with topic:"+topic+";offset:"+offset); } latch.countDown(); } }); } //消息处理完才停止发送者。 latch.await(); producer.close(); } }
  • 消费者消费主流程

    1. 设置Consumer核心属性:可选的属性都可以由ConsumerConfig类管理
    2. 拉取消息:Kafka采用Consumer主动拉取消息的Pull模式,主动从Broker拉取感兴趣的消息
    3. 处理消息,提交位点:消费者拉取完消息后,需要向Broker提交偏移量offset,如果不提交,Broker会认为消费失败,从而再次尝试推送。
    java
    复制代码
    public class MyConsumer { private static final String BOOTSTRAP_SERVERS = "worker1:9092,worker2:9092,worker3:9092"; private static final String TOPIC = "disTopic"; public static void main(String[] args) { //PART1:设置发送者相关属性 Properties props = new Properties(); //kafka地址 props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS); //每个消费者要指定一个group props.put(ConsumerConfig.GROUP_ID_CONFIG, "test"); //key序列化类 props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); //value序列化类 props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); Consumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Arrays.asList(TOPIC)); while (true) { //PART2:拉取消息 // 100毫秒超时时间 ConsumerRecords<String, String> records = consumer.poll(Duration.ofNanos(100)); //PART3:处理消息 for (ConsumerRecord<String, String> record : records) { System.out.println("offset = " + record.offset() + ";key = " + record.key() + "; value= " + record.value()); } //提交offset,消息就不会重复推送。 consumer.commitSync(); //同步提交,表示必须等到offset提交完毕,再去消费下一批数据。 // consumer.commitAsync(); //异步提交,表示发送完提交offset请求后,就开始消费下一批数据了。不用等到Broker的确认。 } } }
  • 工作机制详解

    • 消费者分组消费机制

      在Consumer中,必须要指定一个GROUP_ID_CONFIG属性,表示当前Consumer所属的消费者组。

      image.png

      生产者往topic下发消息时,会尽量均匀的将消息发送到Topic下的各个Partition当中,这个消息会向所有订阅了该Topic的消费者推送,推送时每个ConsumerGroup只会推送一份。同一个消费者组中的多个消费者实例只会共同消费一个消息副本。

      Offset表示每个消费者组在每个Partition中已经消费处理的进度,需要由消费者处理完成后主动向Broker提交,提交完成后Broker才会更新消费进度,如果没有提交,Broker就会认为还没被处理,会往对应的消费者组进行重新投递,重新投递时一般会尽量推送给同一消费者组中的其他消费实例。可将ConsumerConfig中的ENABLE_AUTO_COMMIT_CONFIG参数设置为自动投递。一个Partition最多只能同时被一个Consumer消费。

      但是Offset数据保存在Broker端却由Consumer端维护,显然不能确保数据安全,在ConsumerConfig中,提供了AUTO_OFFSET_RESET_CONFIG参数,来指定消费者组在服务端的Offset不存在时怎么进行消费

      • earliest 自动设置为最早的offset
      • latest 自动设置为最晚的offset
      • none 找不到就抛异常

      这样配置只能算个服务端兜底,消费者要怎么保证offset的安全性?

      • 异步提交:消费者处理业务的同时,异步向Broker提交offset,不阻断业务执行流程,效率高;但是如果消费者消息处理失败,而offset又提交成功,就会造成消息丢失
      • 同步提交:等到业务逻辑执行完成之后再提交,保证休息不丢失;但是同步处理消费速度会变慢,而且有可能产生重复消息,要做好幂等处理
    • 生产者拦截机制

      生产者拦截机制允许客户端在生产者在消息发送到Kafka之前对消息进行拦截,允许有多个拦截类,用逗号隔开即可。

      props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG,"com.roy.kfk.basic.MyInterceptor");

      java
      复制代码
      public class MyInterceptor implements ProducerInterceptor { //发送消息时触发 @Override public ProducerRecord onSend(ProducerRecord producerRecord) { System.out.println("prudocerRecord : " + producerRecord.toString()); return producerRecord; } //收到服务端响应时触发 @Override public void onAcknowledgement(RecordMetadata recordMetadata, Exception e) { System.out.println("acknowledgement recordMetadata:"+recordMetadata.toString()); } //连接关闭时触发 @Override public void close() { System.out.println("producer closed"); } //整理配置项 @Override public void configure(Map<String, ?> map) { System.out.println("=====config start======"); for (Map.Entry<String, ?> entry : map.entrySet()) { System.out.println("entry.key:"+entry.getKey()+" === entry.value: "+entry.getValue()); } System.out.println("=====config end======"); } }
    • 消息序列化机制

      通过KEY_SERIALIZER_CLASS_CONFIG和VALUE_SERIALIZER_CLASS_CONFIG参数指定生产者如何将消息的key和value序列化成二进制数据byte[]数组。

      在Kafka的定义中,key是用来进行分区的可选项,通过key来判断消息发到哪个分区,如果没有填写key,则会自动选择Partition;value是要传递的具体业务消息。

      KEY_DESERIALIZER_CLASS_CONFIG和VALUE_DESERIALIZER_CLASS_CONFIG是指定反序列化参数配置。

      image.png

      实现序列化时,可以简单的转成JSON格式再转二进制数组;也可自定义实现序列化,后者占用的空间更小,序列化也更高效

      • 定长的基础数据类型:直接转成序列化byte[]数组即可
      • 不定长的浮动类型:计算数据长度并作记录,放在数据头部一起序列化,反序列化时按照记录的长度取
    • 消息分区路由机制

      image.png

      • 生产者→Broker

        在ProducerConfig中,PARTITIONER_CLASS_CONFIG参数可以指定实现分区投递的类,默认提供了三种实现类RoundRobinPartitioner,DefaultPartitionerUniformStickyPartitioner,也可以通过实现Partitioner接口的partition方法自定义分配策略。

        生产者默认的Sticky策略在给生产者分配了一个分区后,会尽可能一直使用这个分区。等待该分区的batch.size(默认16k)已满,或者这个分区的消息已完成linger.ms(默认0ms,表示如果batch.size迟迟没有满的等待时间)。RoundRobinPartitioner时在各个Partition中进行轮询发送,没有考虑消息大小及Broker之间的性能差异,较少使用。

      • Broker→消费者

        在ConsumerConfig中,可以指定PARTITION_ASSIGNMENT_STRATEGY分区分配策略,觉得如何在多个Consumer实例和多个Partitioner之间建立关联关系。

        Kafka默认提供了三种消费者的分区分配策略

        • range策略: 比如一个Topic有10个Partiton(partition 0-9) 一个消费者组下有三个Consumer(consumer1-3)。Range策略就会将分区0-3分给一个Consumer,4-6给一个Consumer,7-9给一个Consumer。
        • round-robin策略:轮询分配策略,可以理解为在Consumer中一个一个轮流分配分区。比如0,3,6,9分区给一个Consumer1,1,4,7分区给一个Consumer2,然后2,5,8给一个Consumer3
        • sticky策略:粘性策略。这个策略有两个原则:
          • 1、在开始分区时,尽量保持分区的分配均匀。比如按照Range策略分(这一步实际上是随机的)。
          • 2、分区的分配尽可能的与上一次分配的保持一致。比如在range分区的情况下,第三个Consumer的服务宕机了,那么按照sticky策略,就会保持consumer1和consumer2原有的分区分配情况。然后将consumer3分配的7~9分区尽量平均的分配到另外两个consumer上。这种粘性策略可以很好的保持Consumer的数据稳定性。

        另外可以通过继承AbstractPartitionAssignor抽象类自定义消费者的订阅方式。

    • 生产者消息缓存机制

      Kafka生产者为了避免高并发请求对服务端造成过大压力,每次发消息时并不是一条一条发往服务端,而是增加了一个高速缓存,将消息集中到缓存后,批量进行发送。这种缓存机制也是高并发处理时非常常用的一种机制。

      缓存机制涉及到两个关键组件:accumulator和sender

      java
      复制代码
      //1.记录累加器 int batchSize = Math.max(1, config.getInt(ProducerConfig.BATCH_SIZE_CONFIG)); this.accumulator = new RecordAccumulator(logContext,batchSize,this.compressionType,lingerMs(config),retryBackoffMs,deliveryTimeoutMs, partitionerConfig,metrics,PRODUCER_METRIC_GROUP_NAME,time,apiVersions,transactionManager,new BufferPool(this.totalMemorySize, batchSize, metrics, time, PRODUCER_METRIC_GROUP_NAME)); //2. 数据发送线程 this.sender = newSender(logContext, kafkaClient, this.metadata);

      RecordAccumulator就是Kafka生产者的消息累加器。KafkaProducer要发送的消息都会在RecordAccumulator中缓存起来,然后再分批发送给Kafka broker。在RecordAccumulator中,会针对每一个Partition,维护一个Deque双端队列,这些Dequeue队列基本上是和Kafka服务端的Topic下的Partition对应的。每个Dequeue里会放入若干个ProducerBatch数据。KafkaProducer每次发送的消息,都会根据key分配到对应的Deque队列中。然后每个消息都会保存在这些队列中的某一个ProducerBatch中。再按Partitioner组件的规则完成分发。

      image.png

      Sender就是KafkaProducer中用来发送消息的一个单独的线程。从这里可以看到,每个KafkaProducer对象都对应一个sender线程。他会负责将RecordAccumulator中的消息发送给Kafka。Sender也并不是一次就把RecordAccumulator中缓存的所有消息都发送出去,而是每次只拿一部分消息。只获取RecordAccumulator中缓存内容达到BATCH_SIZE_CONFIG大小的ProducerBatch消息。如果消息比较少,ProducerBatch中的消息大小长期达不到BATCH_SIZE_CONFIG的话,Sender也不会一直等待。最多等待LINGER_MS_CONFIG时长。然后就会将ProducerBatch中的消息读取出来。LINGER_MS_CONFIG默认值是0。

      image.png

      Sender对读取出来的消息,会以Broker为key,缓存到一个对应的队列当中。这些队列当中的消息就称为InflightRequest。接下来这些Inflight就会一一发往Kafka对应的Broker中(Inflight中的消息顺序不保证),直到收到Broker的响应,才会从队列中移除。这些队列也并不会无限缓存,最多缓存MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION(默认值为5)个请求。

    • 发送应答机制

      ProducerConfig下的ACKS_CONFIG配置

      • acks=0,生产者不关心Broker端有没有将消息写入到Partition,只发送消息就不管了。吞吐量是最高的,但是数据安全性是最低的。

      • acks=all or -1,生产者需要等Broker端的所有Partiton(Leader Partition以及其对应的Follower Partition都写完了才能得到返回结果,这样数据是最安全的,但是每次发消息需要等待更长的时间,吞吐量是最低的。适用于敏感数据。

      • acks设置成1,则是一种相对中和的策略。Leader Partition在完成自己的消息写入后,就向生产者返回结果。适用于日志,接受少量数据丢失的场景。

      如果ack设置为all或者-1 ,Kafka也并不是强制要求所有Partition都写入数据后才响应。在Kafka的Broker服务端会有一个配置参数min.insync.replicas,控制Leader Partition在完成多少个Partition的消息写入后,往Producer返回响应。这个参数可以在broker.conf文件中进行配置。

    • 生产者消息幂等性

      当acks设置为1或-1时,每次发送消息都要获取Broker端的RecordMetadata,中间需要经过两次跨网络请求

      image.png

      Producer发送消息的过程中,如果第一步请求成功了, 但是第二步却没有返回。这时,Producer就会认为消息发送失败了。那么Producer必然会发起重试。重试次数由参数ProducerConfig.RETRIES_CONFIG,默认值是Integer.MAX。

      Kafka为了保证消息发送的Exactly-once语义,增加了几个概念:

      • PID:每个新的Producer在初始化的过程中就会被分配一个唯一的PID。这个PID对用户是不可见的。
      • Sequence Numer: 对于每个PID,这个Producer针对Partition会维护一个sequenceNumber。这是一个从0开始单调递增的数字。当Producer要往同一个Partition发送消息时,这个Sequence Number就会加1。然后会随着消息一起发往Broker。
      • Broker端则会针对每个<PID,Partition>维护一个序列号(SN),只有当对应的SequenceNumber = SN+1时,Broker才会接收消息,同时将SN更新为SN+1。否则,SequenceNumber过小就认为消息已经写入了,不需要再重复写入。而如果SequenceNumber过大,就会认为中间可能有数据丢失了。对生产者就会抛出一个OutOfOrderSequenceException。

      Kafka在打开idempotence幂等性控制后,在Broker端就会保证每条消息在一次发送过程中,Broker端最多只会刚刚好持久化一条。这样就能保证at-most-once语义。生产者的acks参数设置成1或-1,保证at-least-once语义,这样就整体上保证了Exactaly-once语义。

      image.png

    • 生产者数据压缩机制

      ProducerConfig中的 COMPRESSION_TYPE_CONFIG配置可以指定选择'gzip', 'snappy', 'lz4', 'zstd’四种压缩算法。

      当生产者往Broker发送消息时,还会对每个消息进行压缩,从而降低Producer到Broker的网络数据传输压力,同时也降低了Broker的数据存储压力,但是压缩也会给CPU带来消耗,根据实际情况选择开启与否。

      在Broker端的broker.conf文件中,也是可以配置压缩算法的。正常情况下,Broker从Producer端接收到消息后不会对其进行任何修改,但是如果Broker端和Producer端指定了不同的压缩算法,就会产生很多异常的表现。

      如果开启了消息压缩,消息从Producer到Broker再到Consumer会一直携带消息的压缩方式,当Consumer读取到消息集合时就知道了这些消息使用的是哪种压缩算法并自行解压。此时要注意的是应用中使用的Kafka客户端版本和Kafka服务端版本是否匹配。

    • 生产者消息事物

      如果生产者一次发送多条消息,然后给不同的消息指定不同的key。这批消息就有可能写入多个Partition,而这些Partition是分布在不同Broker上的。这意味着,Producer需要对多个Broker同时保证消息的幂等性。此时通过上面的生产者消息幂等性机制就无法保证所有消息的幂等了,于是Kafka引入了事物机制,确保同一批消息最好同时成功或失败重试

      java
      复制代码
      // 1 初始化事务 void initTransactions(); // 2 开启事务 void beginTransaction() throws ProducerFencedException; // 3 提交事务 void commitTransaction() throws ProducerFencedException; // 4 放弃事务(类似于回滚事务的操作) void abortTransaction() throws ProducerFencedException;
      java
      复制代码
      public class TransactionErrorDemo { private static final String BOOTSTRAP_SERVERS = "worker1:9092,worker2:9092,worker3:9092"; private static final String TOPIC = "disTopic"; public static void main(String[] args) throws ExecutionException, InterruptedException { Properties props = new Properties(); // 此处配置的是kafka的端口 props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS); // 事务ID props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG,"111"); // 配置key的序列化类 props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer"); // 配置value的序列化类 props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer"); Producer<String,String> producer = new KafkaProducer<>(props); producer.initTransactions(); producer.beginTransaction(); for(int i = 0; i < 5; i++) { ProducerRecord<String, String> record = new ProducerRecord<>(TOPIC, Integer.toString(i), "MyProducer" + i); //异步发送。 producer.send(record); if(i == 3){ //第三条消息放弃事务之后,整个这一批消息都回退了。 System.out.println("error"); producer.abortTransaction(); } } System.out.println("message sended"); try { Thread.sleep(10000); } catch (Exception e) { e.printStackTrace(); } // producer.commitTransaction(); producer.close(); } }

      实际上,Kafka的事务消息还会做两件事情:

      1、一个TransactionId只会对应一个PID

      如果当前一个Producer的事务没有提交,而另一个新的Producer保持相同的TransactionId,这时旧的生产者会立即失效,无法继续发送消息。

      2、跨会话事务对齐

      如果某个Producer实例异常宕机了,事务没有被正常提交。那么新的TransactionId相同的Producer实例会对旧的事务进行补齐。保证旧事务要么提交,要么终止。这样新的Producer实例就可以以一个正常的状态开始工作。

  • SpringBoot 集成

    xml
    复制代码
    <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency>
    xml
    复制代码
    ###########【Kafka集群】########### spring.kafka.bootstrap-servers=worker1:9092,worker2:9093,worker3:9093 ###########【初始化生产者配置】########### # 重试次数 spring.kafka.producer.retries=0 # 应答级别:多少个分区副本备份完成时向生产者发送ack确认(可选0、1、all/-1) spring.kafka.producer.acks=1 # 批量大小 spring.kafka.producer.batch-size=16384 # 提交延时 spring.kafka.producer.properties.linger.ms=0 # 生产端缓冲区大小 spring.kafka.producer.buffer-memory = 33554432 # Kafka提供的序列化和反序列化类 spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer ###########【初始化消费者配置】########### # 默认的消费组ID spring.kafka.consumer.properties.group.id=defaultConsumerGroup # 是否自动提交offset spring.kafka.consumer.enable-auto-commit=true # 提交offset延时(接收到消息后多久提交offset) spring.kafka.consumer.auto-commit-interval=1000 # 当kafka中没有初始offset或offset超出范围时将自动重置offset # earliest:重置为分区中最小的offset; # latest:重置为分区中最新的offset(消费分区中新产生的数据); # none:只要有一个分区不存在已提交的offset,就抛出异常; spring.kafka.consumer.auto-offset-reset=latest # 消费会话超时时间(超过这个时间consumer没有发送心跳,就会触发rebalance操作) spring.kafka.consumer.properties.session.timeout.ms=120000 # 消费请求超时时间 spring.kafka.consumer.properties.request.timeout.ms=180000 # Kafka提供的序列化和反序列化类 spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer
    java
    复制代码
    // 生产者 @RestController public class KafkaProducer { @Autowired private KafkaTemplate<String, Object> kafkaTemplate; // 发送消息 @GetMapping("/kafka/normal/{message}") public void sendMessage1(@PathVariable("message") String normalMessage) { kafkaTemplate.send("topic1", normalMessage); } } // 消费者 @Component public class KafkaConsumer { // 消费监听 @KafkaListener(topics = {"topic1"}) public void onMessage1(ConsumerRecord<?, ?> record){ // 消费的哪个topic、partition的消息,打印出消息内容 System.out.println("简单消费:"+record.topic()+"-"+record.partition()+"-"+record.value()); } }
0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP