编程导航Kafka话题讨论

Kafka

5 参与
分享

快来分享你的内容吧~

点击登录,快来和大家讨论吧~
表情
图片
话题
打卡
综合
交流
文章
问答

Day26 MQ补充Kafka相关

### 客户端消息流转流程 ![image.png](https://pic.code-nav.cn/post_picture/1876274222060195841/fL1YHI7NwmWDaOPx.webp) - 生产者发送主流程 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](https://pic.code-nav.cn/post_picture/1876274222060195841/R3dwMuRCNYWJGe3d.webp) 生产者往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](https://pic.code-nav.cn/post_picture/1876274222060195841/EkagnahaDWn2VrYk.webp) 实现序列化时,可以简单的转成JSON格式再转二进制数组;也可自定义实现序列化,后者占用的空间更小,序列化也更高效 - 定长的基础数据类型:直接转成序列化byte[]数组即可 - 不定长的浮动类型:计算数据长度并作记录,放在数据头部一起序列化,反序列化时按照记录的长度取 - 消息分区路由机制 ![image.png](https://pic.code-nav.cn/post_picture/1876274222060195841/nJ1dobWsRV3LzPdV.webp) - **生产者→Broker** 在ProducerConfig中,**PARTITIONER_CLASS_CONFIG**参数可以指定实现分区投递的类,默认提供了三种实现类RoundRobinPartitioner,~~DefaultPartitioner~~和~~UniformStickyPartitioner~~,也可以通过实现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](https://pic.code-nav.cn/post_picture/1876274222060195841/Iyp96NqYJbkBYxbP.webp) **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](https://pic.code-nav.cn/post_picture/1876274222060195841/mTzqNhav0EsI8xlJ.webp) 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](https://pic.code-nav.cn/post_picture/1876274222060195841/BWarEm983nCootfi.webp) 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](https://pic.code-nav.cn/post_picture/1876274222060195841/PdCiLf3Ggw9lFmi7.webp) - 生产者数据压缩机制 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()); } } ```

Kafka重平衡机制

## Kafka架构 消费者 ## 重平衡 重平衡是指在消费者加入或离开消费者群组时,由消费者协调器(Coordinator)发起的重新分配分区的过程。在重平衡过程中,消费者会停止读取消息,释放已经持有的分区并重新分配新的分区,从而实现消费者负载均衡,避免某些消费者处理过多的消息,而其他消费者处于空闲状态。 重平衡只是影响消费者的分区,即在消费者发生变动的情况下进行更改的(挂了什么的) ## 重平衡策略 Kafka是一款开源的分布式消息队列,支持多个消费者同时订阅同一个topic。为了保证消费者集群内各个节点的负载均衡,Kafka提供了四种重平衡策略。 范围分配(平均划定范围) 轮询分区(默认的策略) 模板匹配 粘性分区(在进行新的分配之前考虑上一次分配的结果,减少重分配) ## 重平衡如何进行 ### 进行的时机 当 Kafka 集群中添加或删除主题分区时,或者消费者加入或离开消费组时,就会发生重平衡。重平衡会导致消费者重新分配分区,这意味着该消费者可能需要重新加载从其他消费者分配来的分区数据。 ### 从消费者端看重平衡 ### (谁先加入谁是主要消费者负责人)-主要负责消费组中消费者分区的分配(消费者协调器) ### 步骤(发送两次消息joingroup和SyncGroup请求) 1,先确定人数 2,再确定策略并且将策略进行广播 ## 从协调者(组协调器)的角度看(消费者进组) 分为很多种情况,但都是基于心跳检测进行的 ### 消费者组的5种状态 ​ 5种状态分别为 ![新的消费者从empty开始](https://pic.code-nav.cn/post_picture/1829067802410602497/RgKQrpQ9p5oy2dli.webp) ## 其中的重要参数 session.timeout.ms(检测消费者失败的时间),更小则更快发现重平衡避免消费滞后,但是也会导致频繁重平衡 max.poll.interval.ms(消费者处理消息逻辑的最大时间,对于某些业务来说,处理消息可能需要很长时间,比如需要 1分钟,那么该参数就需要设置成大于 1分钟的值,否则就会被 Coordinator(协调者) 剔除消息组然后重平衡。) heartbeat.interval.ms(该参数跟 session.timeout.ms 紧密关联,只要在 session.timeout.ms 时间内与Coordinator 保持心跳,就不会被 Coordinator 剔除,那么心跳间隔的时间就是session.timeout.ms,因此,该参数值必须小于 session.timeout.ms ,以保持session.timeout.ms 时间内有心跳。)

Kafka偏移量

## 生产者偏移量 将生产消息写入kafka中,kafka为该消息分配的在分区中的顺序编号 ## 作用: **消息的唯一标识**: 每条消息在其所属的分区中都有一个唯一的偏移量,生产者的偏移量通常由 Kafka 自动生成和分配,用于确保消息在分区中的顺序性和唯一性。 **消费消息的起点**: 生产者的偏移量影响消息在分区中的位置,消费者依赖这些偏移量来有序地读取消息。消费者的消费过程总是根据生产者生成的偏移量逐条读取消息 ## 消费者偏移量(意义-提交方式) 消费者偏移量是指消费者在 Kafka 中消费到某个主题分区时,标识其当前消费位置的标记。 ## 作用: **消息跟踪**: 通过消费者偏移量,Kafka 能知道消费者已经消费到哪一条消息,便于消费者在断开连接或重启时,继续从上次消费的地方继续消费,保证消息不重复消费或丢失。 **消费管理**: Kafka 允许消费者手动或自动提交偏移量,以确保在消费过程中记录消费进度。消费者可以选择手动控制偏移量提交时间,或让 Kafka 自动完成。 **实现消息重新处理**: 消费者可以根据需要重置偏移量,从而重新消费某些消息。例如,在处理逻辑失败后,可以将偏移量重置到之前的位置重新消费消息。 ## 存储(持久化和容灾) 消费者偏移量通常会被存储在 Kafka 内部的特殊主题(默认是 `__consumer_offsets`),也可以存储在外部如 Zookeeper 或自定义数据库中 ## 提交方式 ### 自动提交 缺点: 自动提交间隔时间会导致重复消费,比如自动提交时间100S,首次提交的偏移量是20,而消费者拉取了5条消息消费了 ,在50秒的时候broker突然宕机,发生了分区再均衡,这个时候之前分区的消息从之前的消费者转移到另外一个新消费者,新的消费者会读取50秒之前提交的偏移量,导致重复消费。 丢失消息 消费者批量推送消息时,消费之只消费了20条,如果broker宕机,发生分区再均衡时,会从上次提交的偏移量位置重新消费,导致批量的消息(80条)丢失 注意:在调用close方法之前,是会自动提交一次偏移量,但异常或者提前退出轮询(poll)是不会的,需要自己定义策略来保证。比如finally方法手动提交commitSync() ### 手动提交 ### 异步提交 ### 异步+同步提交 ### 提交特定的偏移量

Kafka知识点总结(Java版本)

# Kafka:分布式流处理平台的关键概念与技术特性 Kafka,作为一款高性能的分布式流处理平台,以其卓越的数据处理能力和灵活的架构设计,在实时数据管道和应用程序构建中扮演着重要角色。以下是Kafka的核心知识点,旨在为您揭示其背后的技术魅力。 ## 一、基础架构与核心术语 生产者(Producer):作为消息的发起者,生产者负责将数据推送至Kafka集群,是数据流转的起点。 消费者(Consumer):作为消息的接收者,消费者从Kafka集群中提取数据,实现信息的最终消费。 代理(Broker):Kafka集群中的服务器节点,它们存储数据并处理客户端的请求,是整个系统的骨架。 主题(Topic):Kafka中消息的逻辑分类,类似于数据库中的表,用于区分不同类型的数据。 分区(Partition):主题的物理分组,每个分区都是一个有序且不可变的消息序列,增强了Kafka的并发处理能力。 偏移量(Offset):分区中消息的唯一标识,用于定位消息在分区中的位置。 ## 二、架构优势与数据持久化 Kafka的设计理念旨在提供以下架构优势: 分布式部署:通过多Broker节点组成集群,实现数据的分布式存储与处理。 高吞吐量:利用分区机制和批量发送技术,Kafka能够实现高效的数据传输。 可扩展性:随着业务需求的增长,Kafka可通过增加Broker节点轻松扩展集群规模。 数据持久化:Kafka将消息持久化到磁盘,确保数据的长期保存,防止数据丢失。 ## 三、副本机制与容错能力 Kafka采用副本机制,为数据提供冗余保护,提升系统的容错能力: 每个分区都有多个副本,其中一个是Leader副本,负责处理读写请求,其他为Follower副本。 当Leader副本发生故障时,Kafka能够自动进行故障转移,由其他Follower副本接替领导角色。 ## 四、消费者组与消息消费 Kafka中的消费者组机制允许多个消费者协同工作,共同消费一个主题的消息: 每个消费者组独立消费主题中的消息,组内消费者之间不会重复消费。 通过消费者组,Kafka实现了消息的负载均衡和并行处理。 ## 五、生产者与消费者的交互 生产者负责将消息发送到指定的主题,支持同步和异步发送模式。 消费者可以根据需要按顺序消费消息,或通过指定偏移量进行任意位置的消息消费。 ## 六、高级特性与应用 Kafka的高级特性进一步提升了其数据处理能力: 消息压缩:支持多种压缩算法,减少网络传输数据量,提高传输效率。 消息事务:确保消息的ACID特性,满足事务性消息处理的需求。 流处理:提供流处理API,允许用户构建实时数据流处理应用,实现数据的即时转换和分析。 ### Kafka生产者示例代码 ```java import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; public class KafkaProducerExample { public static void main(String[] args) { // 设置Kafka生产者配置 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); // Kafka集群地址 props.put("key.serializer", StringSerializer.class.getName()); props.put("value.serializer", StringSerializer.class.getName()); // 创建Kafka生产者 KafkaProducer<String, String> producer = new KafkaProducer<>(props); // 构建消息 String topic = "test-topic"; String key = "key1"; String value = "Hello, Kafka!"; // 发送消息 ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value); producer.send(record); // 关闭生产者 producer.close(); } } ``` ### Kafka消费者示例代码 ```java import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class KafkaConsumerExample { public static void main(String[] args) { // 设置Kafka消费者配置 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); // Kafka集群地址 props.put("group.id", "test-group"); // 消费者组ID props.put("key.deserializer", StringDeserializer.class.getName()); props.put("value.deserializer", StringDeserializer.class.getName()); props.put("auto.offset.reset", "earliest"); // 从最早的offset开始消费 // 创建Kafka消费者 KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); // 订阅主题 String topic = "test-topic"; consumer.subscribe(Collections.singletonList(topic)); // 拉取消息并消费 while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, String> record : records) { System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value()); } } // 实际应用中需要正确处理消费者的关闭操作,通常在接收到关闭信号时关闭消费者 consumer.close(); } } ``` ## 数据一致性 Kafka的设计目标是:高吞吐、高并发、高性能。为了做到以上三点,它必须设计成分布式的,多台机器可以同时提供读写,并且需要为数据的存储做冗余备份。 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/aVNDOtEe-image.webp" alt="image.png" width="100%" /> 图中的主题有3个分区,每个分区有3个副本,这样数据可以冗余存储,提高了数据的可用性。并且3个副本有两种角色,Leader和Follower,Follower副本会同步Leader副本的数据。 一旦Leader副本挂了,Follower副本可以选举成为新的Leader副本, 这样就提升了分区可用性,但是相对的,在提升了分区可用性的同时,也就牺牲了数据的一致性。 我们来看这样的一个场景:一个分区有3个副本,一个Leader和两个Follower。Leader副本作为数据的读写副本,所以生产者的数据都会发送给leader副本,而两个follower副本会周期性地同步leader副本的数据,但是因为网络,资源等因素的制约,同步数据的过程是有一定延迟的,所以3个副本之间的数据可能是不同的。具体如下图所示: <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/Kmksvb8m-image.webp" alt="image.png" width="339px" /> 此时,假设leader副本因为意外原因宕掉了,那么Kafka为了提高分区可用性,此时会选择2个follower副本中的一个作为Leader对外提供数据服务。此时我们就会发现,对于消费者而言,之前leader副本能访问的数据是D,但是重新选择leader副本后,能访问的数据就变成了C,这样消费者就会认为数据丢失了,也就是所谓的数据不一致了。 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/jAdDxWdz-image.webp" alt="image.png" width="245px" /> 为了提升数据的一致性,Kafka引入了高水位(HW :High Watermark)机制,Kafka在不同的副本之间维护了一个水位线的机制(其实也是一个偏移量的概念),消费者只能读取到水位线以下的的数据。这就是所谓的木桶理论:木桶中容纳水的高度,只能是水桶中最短的那块木板的高度。这里将整个分区看成一个木桶,其中的数据看成水,而每一个副本就是木桶上的一块木板,那么这个分区(木桶)可以被消费者消费的数据(容纳的水)其实就是数据最少的那个副本的最后数据位置(木板高度)。 也就是说,消费者一开始在消费Leader的时候,虽然Leader副本中已经有a、b、c、d 4条数据,但是由于高水位线的限制,所以也只能消费到a、b这两条数据。 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/tlTgAjkW-image.webp" alt="image.png" width="393px" /> 这样即使leader挂掉了,但是对于消费者来讲,消费到的数据其实还是一样的,因为它能看到的数据是一样的,也就是说,消费者不会认为数据不一致。 不过也要注意,因为follower要求和leader的日志数据严格保持一致,所以就需要根据现在Leader的数据偏移量值对其他的副本进行数据截断(truncate)操作。 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/sFDBickQ-image.webp" alt="image.png" width="392px" />

Kafka集群Windows配置

# 配置成功的效果 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/1nxq4aih-image.webp" alt="image.png" width="100%" /> ## 详细配置 ### 环境安装 当前Java软件开发中,主流的版本就是Java 8,而Kafka 3.X官方建议Java版本更新至Java11,但是Java8依然可用。未来Kafka 4.X版本会完全弃用Java8,不过,咱们当前学习的Kafka版本为3.6.1版本,所以使用Java8即可,无需升级。也可以使用Java8以上的版本 安装Kafka 下载软件安装包:kafka_2.12-3.6.1.tgz,[下载地址](https://kafka.apache.org/downloads) 这里演示版本是3.6.1,是Kafka软件的版本。 2.12是对应的Scala开发语言版本。Scala2.12和Java8是兼容的,所以可以直接使用。 tgz是一种linux系统中常见的压缩文件格式,类似与windows系统的zip和rar格式。所以Windows环境中可以直接使用压缩工具进行解压缩。 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/9Qlu9nox-image.webp" alt="image.png" width="100%" /> ## 解压文件 kafka_2.12-3.6.1.tgz,解压目录为非系统盘的根目录,比如e:/ <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/I5DsqhLr-image.webp" alt="image.png" width="100%" /> 为了访问方便,可以将解压后的文件目录改为kafka, 更改后的文件目录结构如下: | 文件夹名称 |功能 | | --- | --- | | bin | linux系统下可执行脚本文件 | | bin/windows | windows系统下可执行脚本文件 | |config|配置文件| |libs|依赖类库| |licenses|许可信息| |site-docs|文档| |logs|服务日志| (1)在磁盘根目录创建文件夹cluster,文件夹名称不要太长 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/NVO3Q7dn-image.png" alt="image.png" width="100%" /> (2)将kafka安装包kafka-3.6.1-src.tgz解压缩到kafka文件夹 ### 安装ZooKeeper (1)修改文件夹名为kafka-zookeeper 因为kafka内置了ZooKeeper软件,所以此处将解压缩的文件作为ZooKeeper软件使用。 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/imtlM0EC-image.webp" alt="image.png" width="282px" /> (2)修改config/zookeeper.properties文件 ``` # Licensed to the Apache Software Foundation (ASF) under one or more # contributor license agreements. See the NOTICE file distributed with # this work for additional information regarding copyright ownership. # The ASF licenses this file to You under the Apache License, Version 2.0 # (the "License"); you may not use this file except in compliance with # the License. You may obtain a copy of the License at # # http://www.apache.org/licenses/LICENSE-2.0 # # Unless required by applicable law or agreed to in writing, software # distributed under the License is distributed on an "AS IS" BASIS, # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. # the directory where the snapshot is stored. # 此处注意,如果文件目录不存在,会自动创建 dataDir=E:/cluster/kafka-zookeeper/data # the port at which the clients will connect # ZooKeeper默认端口为2181 clientPort=2181 # disable the per-ip limit on the number of connections since this is a non-production config maxClientCnxns=0 # Disable the adminserver by default to avoid port conflicts. # Set the port to something non-conflicting if choosing to enable this admin.enableServer=false # admin.serverPort=8080 ``` ### 安装Kafka (1)将上面解压缩的文件复制一份,改名为kafka_broker_1 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/IfTzMYio-image.png" alt="image.png" width="100%" /> (2)修改config/server.properties配置文件 ```bash # Licensed to the Apache Software Foundation (ASF) under one or more # contributor license agreements. See the NOTICE file distributed with # this work for additional information regarding copyright ownership. # The ASF licenses this file to You under the Apache License, Version 2.0 # (the "License"); you may not use this file except in compliance with # the License. You may obtain a copy of the License at # # http://www.apache.org/licenses/LICENSE-2.0 # # Unless required by applicable law or agreed to in writing, software # distributed under the License is distributed on an "AS IS" BASIS, # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. # # This configuration file is intended for use in ZK-based mode, where Apache ZooKeeper is required. # See kafka.server.KafkaConfig for additional details and defaults # ############################# Server Basics ############################# # The id of the broker. This must be set to a unique integer for each broker. # kafka节点数字标识,集群内具有唯一性 broker.id=1 ############################# Socket Server Settings ############################# # The address the socket server listens on. If not configured, the host name will be equal to the value of # java.net.InetAddress.getCanonicalHostName(), with PLAINTEXT listener name, and port 9092. # FORMAT: # listeners = listener_name://host_name:port # EXAMPLE: # listeners = PLAINTEXT://your.host.name:9092 # 监听器 9091为本地端口,如果冲突,请重新指定 listeners=PLAINTEXT://:9091 # Listener name, hostname and port the broker will advertise to clients. # If not set, it uses the value for "listeners". #advertised.listeners=PLAINTEXT://:9091 # Maps listener names to security protocols, the default is for them to be the same. See the config documentation for more details #listener.security.protocol.map=PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL # The number of threads that the server uses for receiving requests from the network and sending responses to the network num.network.threads=3 # The number of threads that the server uses for processing requests, which may include disk I/O num.io.threads=8 # The send buffer (SO_SNDBUF) used by the socket server socket.send.buffer.bytes=102400 # The receive buffer (SO_RCVBUF) used by the socket server socket.receive.buffer.bytes=102400 # The maximum size of a request that the socket server will accept (protection against OOM) socket.request.max.bytes=104857600 ############################# Log Basics ############################# # A comma separated list of directories under which to store log files # 数据文件路径,如果不存在,会自动创建 log.dirs=E:/cluster/kafka-node-1/data # The default number of log partitions per topic. More partitions allow greater # parallelism for consumption, but this will also result in more files across # the brokers. num.partitions=1 # The number of threads per data directory to be used for log recovery at startup and flushing at shutdown. # This value is recommended to be increased for installations with data dirs located in RAID array. num.recovery.threads.per.data.dir=1 ############################# Internal Topic Settings ############################# # The replication factor for the group metadata internal topics "__consumer_offsets" and "__transaction_state" # For anything other than development testing, a value greater than 1 is recommended to ensure availability such as 3. offsets.topic.replication.factor=1 transaction.state.log.replication.factor=1 transaction.state.log.min.isr=1 ############################# Log Flush Policy ############################# # Messages are immediately written to the filesystem but by default we only fsync() to sync # the OS cache lazily. The following configurations control the flush of data to disk. # There are a few important trade-offs here: # 1. Durability: Unflushed data may be lost if you are not using replication. # 2. Latency: Very large flush intervals may lead to latency spikes when the flush does occur as there will be a lot of data to flush. # 3. Throughput: The flush is generally the most expensive operation, and a small flush interval may lead to excessive seeks. # The settings below allow one to configure the flush policy to flush data after a period of time or # every N messages (or both). This can be done globally and overridden on a per-topic basis. # The number of messages to accept before forcing a flush of data to disk #log.flush.interval.messages=10000 # The maximum amount of time a message can sit in a log before we force a flush #log.flush.interval.ms=1000 ############################# Log Retention Policy ############################# # The following configurations control the disposal of log segments. The policy can # be set to delete segments after a period of time, or after a given size has accumulated. # A segment will be deleted whenever *either* of these criteria are met. Deletion always happens # from the end of the log. # The minimum age of a log file to be eligible for deletion due to age log.retention.hours=168 # A size-based retention policy for logs. Segments are pruned from the log unless the remaining # segments drop below log.retention.bytes. Functions independently of log.retention.hours. #log.retention.bytes=1073741824 # The maximum size of a log segment file. When this size is reached a new log segment will be created. #log.segment.bytes=1073741824 log.segment.bytes=190 log.flush.interval.messages=2 log.index.interval.bytes=17 # The interval at which log segments are checked to see if they can be deleted according # to the retention policies log.retention.check.interval.ms=300000 ############################# Zookeeper ############################# # Zookeeper connection string (see zookeeper docs for details). # This is a comma separated host:port pairs, each corresponding to a zk # server. e.g. "127.0.0.1:3000,127.0.0.1:3001,127.0.0.1:3002". # You can also append an optional chroot string to the urls to specify the # root directory for all kafka znodes. # ZooKeeper软件连接地址,2181为默认的ZK端口号 /kafka 为ZK的管理节点 zookeeper.connect=localhost:2181/kafka # Timeout in ms for connecting to zookeeper zookeeper.connection.timeout.ms=18000 ############################# Group Coordinator Settings ############################# # The following configuration specifies the time, in milliseconds, that the GroupCoordinator will delay the initial consumer rebalance. # The rebalance will be further delayed by the value of group.initial.rebalance.delay.ms as new members join the group, up to a maximum of max.poll.interval.ms. # The default value for this is 3 seconds. # We override this to 0 here as it makes for a better out-of-the-box experience for development and testing. # However, in production environments the default value of 3 seconds is more suitable as this will help to avoid unnecessary, and potentially expensive, rebalances during application startup. group.initial.rebalance.delay.ms=0 ``` (3)将kafka_broker_1文件夹复制两份,改名为kafka_broker_2,kafka_broker_3 (4)分别修改kafka-node-2,kafka-node-3文件夹中的配置文件server.properties 将文件内容中的broker.id=1分别改为broker.id=2,broker.id=3 将文件内容中的9091分别改为9092,9093(如果端口冲突,请重新设置) 将文件内容中的kafka-node-1分别改为kafka-node-2,kafka-node-3 ## 封装启动脚本 因为Kafka启动前,必须先启动ZooKeeper,并且Kafka集群中有多个节点需要启动,所以启动过程比较繁琐,这里我们将启动的指令进行封装。 (1)在kafka-zookeeper文件夹下创建zk.cmd批处理文件 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/cEQgroV0-image.png" alt="image.png" width="100%" /> (2)在zk.cmd文件中添加内容 ```# 添加启动命令 call bin/windows/zookeeper-server-start.bat config/zookeeper.properties ``` (3)在kafka_broker_1,kafka_broker_2,kafka_broker_3文件夹下分别创建kfk.cmd批处理文件 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/Q4Z57hZQ-image.png" alt="image.png" width="100%" /> ```# 添加启动命令 call bin/windows/kafka-server-start.bat config/server.properties ``` (5)在cluster文件夹下创建cluster.cmd批处理文件,用于启动kafka集群 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/cnso0hf9-image.png" alt="image.png" width="100%" /> (6)在cluster.cmd文件中添加内容 ``` cd kafka_zookeeper start zk.cmd ping 127.0.0.1 -n 10 >nul cd ../kafka_broker_1 start kfk.cmd cd ../kafka_broker_2 start kfk.cmd cd ../kafka_broker_3 start kfk.cmd ``` (7)在cluster文件夹下创建cluster-clear.cmd批处理文件,用于清理和重置kafka数据 ``` cd kafka_zookeeper rd /s /q data cd ../kafka_broker_1 rd /s /q data cd ../kafka_broker_2 rd /s /q data cd ../kafka_broker_3 rd /s /q data ``` (9)双击执行cluster.cmd文件,启动Kafka集群 集群启动命令后,会打开多个黑窗口,每一个窗口都是一个kafka服务,请不要关闭,一旦关闭,对应的kafka服务就停止了。如果启动过程报错,主要是因为zookeeper和kafka的同步问题,请先执行cluster-clear.cmd文件,再执行cluster.cmd文件即可。 <img src="https://pic.code-nav.cn/post_picture/1623004460282298369/eJUDol8u-7a56fddbc5a334f2c18861f27e1f743.webp" alt="7a56fddbc5a334f2c18861f27e1f743.png" width="100%" />

下载 APP