Kafka 入门学习
目录1 初识KafKa1.1基本概念2 生产者2.1客户端开发2.1.1 必要参数2.1.2 消息的发送2.1.3 序列化2.1.4 分区器2.1.5 生产者拦截器2.2 整体架构2.2.1 RecordAccumulator2.2.2 Sender线程3 消费者3.1 消费者和消费者组3.1.1 消息投递模式3.2 客户端开发3.2.1 必要参数3.2.2 订阅主题与分区3.2.3 反序列化3.2.4 消费消息3.2.5 位移提交3.2.6 控制或关闭消费3.2.7 指定位移消费3.2.8 再均衡3.2.9 消费者拦截器3.2.10 多线程实现3.2.11 重要消费参数4 主题与分区4.1 优先副本的选举4.2 文件目录4.3 日志索引4.4 日志快速读取5、队列1 初识KafKa1.1基本概念1.Producer:生产者生产者负责创建消息投递到Kafka中2.Consumer:消费者连接到KafKa上并接收消息进行处理。3.Broker独立的Kafka服务节点或者服务实例。4.Topic: Kafka中消息以主题为单位进行归类。生产这将消息发送到特定的主题每一个消息都需要指定主题消费者负责订阅主题并进行消费。5.Partition:主题是一个逻辑上的概念他可以细分为多个分区。同一主题下的不同分区包含的消息不同。分区在存储层面可以看做是一个追加的日志文件。消息被追加到日志文件会分配一个特定的偏移量offest,offest是分区中的唯一标识offest不会跨越分区所以只保证分区中的消息有序。分区可以分布在不通的服务器上也就是说一个主题可以横跨多个broker。可以解决单文件只能在一个服务器上造成的性能问题。6.Replica:Kafka为分区引入了多副本概念增加分区副本数量可以提升容灾能力。副本同一时间并非完全一样一主多从leader副本负责读写。follower副本只负责消息同步。副本处于不通broker中。当leader出现故障从follower中重新选举新的leader。7.分区中的所有副本(leaderfollower)统称为AR(Assigned Replicas),所有与leader副本保持一定程度同步的副本包括leader组成ISRIn-Sync Replicas.与leader副本同步滞后过多的副本不包过leader组成OSR(out-of-Syn Replicas).leader副本负责维护和跟踪ISR集合中所有的follower副本的滞后状态当follower滞后太多或者失效时leader将其从ISR中剔除。如果OSR中有follower副本追上那么从OSR转移到ISR. 默认情况下当leader发生故障只有ISR中的副本才有资格被宣威leader。ISR与HW和LEO也有密切的关系。HW-Hight Watermark的缩写。高水位。他表示了一个特定的消息偏移量offest,消费者只能拉取到这个offest之前的消息。LEO为Log End Offest缩写。分区中当前日志文件一条待写入消息的offest分区中消息是从Log Start Offest为0开始到LogEndOffest结束。HW就是所有ISR集合中LogEndOffest的最小值2 生产者2.1客户端开发一个正常的生产逻辑需要具备以下几个步骤配置生产者客户端参数以及创建相应的生产者实例。构建待发送的消息发送消息关闭生产者实例//生产者实例 是线程安全的 KafkaProducerString,String prodcuer new KafkaProducer(propos); 发送的消息类。 public class ProducerRecordK, V { //主题 private final String topic; //分区号 private final Integer partition; //消息头 private final Headers headers; //消息key private final K key; //消息值 private final V value; //消息的时间戳 private final Long timestamp;2.1.1 必要参数bootstrap.servers 用来指定连接Kafka集群所需的broker地址清单多个用逗号分割。key.serializervalue.serializer :borker端接收的消息必须以字节数组形式存在。2.1.2 消息的发送发送消息有三种方式fire-and-forget 发后即忘 只管发送不管是否到达sync 同步async异步try{ FutureRecordMetedata future producer.send(record); RecordMetedata metedata future.get(); //可以通过get方法来阻塞等待Kafka的响应直到消息发送成功 }catch(ExecutionException | InterruptedException e){ e.printStackTrace(); }2.1.3 序列化生产者使用序列化把对象转为字节数组才能发送给kafka, 消费者需要用对应的反序列化将字节数组转化为对象。2.1.4 分区器消息在通过send方法发往broker过程中有可能需要经过拦截器Interceptor、序列化器Serializer和分区器Partitioner的一系列作用之后才能被真正的发往broker。作用为消息分配分区。没有指定分区的时候分区器会指定一个默认分区器是org.apache.kafka.clients.producer.internals.DefaultPartitioner。其中partition用来计算分区号返回值为int类型。Partitioner是DefaultPartitioner的父类接口继承了Configurable接口通过该接口中的configure方法获取配置信息。/** * Compute the partition for the given record. * * param topic The topic name * param numPartitions The number of partitions of the given {code topic} * param key The key to partition on (or null if no key) * param keyBytes serialized key to partition on (or null if no key) * param value The value to partition on or null * param valueBytes serialized value to partition on or null * param cluster The current cluster metadata */ public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster, int numPartitions) { if (keyBytes null) { return stickyPartitionCache.partition(topic, cluster); } // hash the keyBytes to choose a partition return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions; }如果key不为null,那么对key进行hash然后计算分区号。相同的key写入同一个分区。如果key为null,那么消息会随机的方式发往主题内的任意一个分区。2.1.5 生产者拦截器可以在发送消息前做一些准备工作比如过滤不合要求的消息修改消息内容。也可以做一些定制化的需求比如统计类工作。使用自定义实现ProducerInterceptor接口。需要实现3个接口public interface ProducerInterceptorK, V extends Configurable { public ProducerRecordK, V onSend(ProducerRecordK, V record); public void onAcknowledgement(RecordMetadata metadata, Exception exception); public void close(); }KafkaProducer在将消息序列化和计算分区前会调用生产者拦截器的onSend方法来对消息进行相应的定制化操作。KafkaProducer会在消息被应答Acknowledge之前或者消息发送失败时调用onAcknowledgement方法优先与用户设定的CallBack之前执行。2.2 整体架构整个生产者客户端由两个线程协调运行这2个线程分别为主线程和Sender线程。主线程中由KafkaProducer创建消息通过可能的拦截器、序列化器、分区器的作用后缓存消息到消息累加器RecordAccumulator。Sender线程从RecordAccumulator中获取消息并将期发送到Kafka中。sender线程也是在构造函数里启动的this.sender newSender(logContext, kafkaClient, this.metadata); String ioThreadName NETWORK_THREAD_PREFIX | clientId; this.ioThread new KafkaThread(ioThreadName, this.sender, true); this.ioThread.start();2.2.1 RecordAccumulator//在KafkaProducer的构造方法中初始化 this.accumulator new RecordAccumulator(logContext, config.getInt(ProducerConfig.BATCH_SIZE_CONFIG), //batchSiz 默认16384 this.compressionType, //消息压缩方式none不压缩、gzip、snappy、lz4、zstd、 lingerMs(config), //用来 retryBackoffMs, deliveryTimeoutMs, metrics, PRODUCER_METRIC_GROUP_NAME, time, apiVersions, transactionManager, new BufferPool(this.totalMemorySize, config.getInt(ProducerConfig.BATCH_SIZE_CONFIG), metrics, time, PRODUCER_METRIC_GROUP_NAME));batchSize :初始化MemoryRecords实例时分配的大小CompressionType:压缩消息的方式默认none不压缩lingerMs, 用来指定生产者发送ProducerBatch之前需要等待更多消息ProducerRecord加入ProducerBatch的时间。默认为0。retryBackoffMs配置生产者重试的次数默认为0。异常情况下会重试。deliveryTimeoutMsrequestTime(Producer请求等待响应的最长时间) lingerMs deliveryTimeoutMs(传递超时时间)metrics:transactionManager:事务管理BufferPool字节缓存池属性ConcurrentMapTopicPartition, DequeProducerBatch batches;以TopicPartition为key, 双端队列为value。 收集器可以收集不同分区的消息各自分区下有一个队列队列中有多个ProducerBatch每个ProducerBatch中可以放很多消息。发送消息追加后还会在IncompleteBatches中缓存起来直到被ACK确认。2.2.2 Sender线程sender线程中维护了确认机制acks 1 发送消息leader副本写入消息就会响应成功。默认0生产者发送消息不需要响应。-1或null发送消息需要ISR都写入成功才响应成功。发送方式每次线程启动只发送一次。 Producer new完之后发送消息完成需要调用close方法会关闭sender线程。3 消费者3.1 消费者和消费者组消费者Consumer负责订阅Kafka中的主题Topic,并从订阅的主题上拉取消息。消费者组Consumer Group每个消费者都有一个消费者组消息发布到主题后只会被投递给订阅他的每个消费者组中的其中一个消费者。注意同一个topic下会有很多的分区同一个消费组中可以增加消费者来让消费能力提升。但是当消费者过多就会导致有的消费者分配不到任何分区。分配策略是通过消费这客户端参数partition.assignment.strategy来配置的。3.1.1 消息投递模式点对点p2p 如果所有消费者都属于同一个消费组。那么所有消息就会被均衡的投递到每一个消费者。每条消息只会被一个消费者消费。发布订阅如果所有消费者都属于不同的消费组。那么所有的消息都会被广播给所有的消费者。那么一个消息会被所有消费者处理。3.2 客户端开发一个正常的客户端需要一下几个步骤配置消费者客户端参数及创建相应的消费者实例。订阅主题拉取消息提交消费位移关闭消费者实例。3.2.1 必要参数bootstrap.servers。 和生产者中一样group.id: 消费者组的名称默认为“”。如果为空会报错。一般会设置成具有一定业务意义的名称。key.deserializer 和 value.deserializer:用于反序列化。3.2.2 订阅主题与分区订阅方式AUTO_TOPICS:集合订阅的方式 AUTO_PATTERN:正则订阅方式 USER_ASSIGNED:assign方式三种方式互斥一个消费者只能使用一种。否则会报错。1.消费者可以订阅一个或者多个主题。consumer.subscribe(Arrays.asList(topic1)); consumer.subscribe(Arrays.asList(topic2)); consumer.subscribe(Pattern.complie(topic-.*));可以使用集合或者正则表达式的形式订阅特定模式的主题。如果前后2次订阅了不通的主题以最后一次为准。2.消费者还可以通过KafkaConsumer中的assign()方法订阅主题主题中特定的分区。public void assign(CollectionTopicPartition partitions) //通过该方法可以获取主题下的所有分区信息 包括 AR ISR OSR集合 public ListPartitionInfo partitionsFor(String topic, Duration timeout)有以上方法所以我们可以通过assign方法也能实现订阅主题全部分区的功能。3.消费者可以通过unsubscribe()方法来取消主题的订阅。3.2.3 反序列化生产者使用序列化把对象转为字节数组才能发送给kafka, 消费者需要用对应的反序列化将字节数组转化为对象。都可以自定义序列化方式不过生产者和消费者得配对。3.2.4 消费消息Kafka中的消费是基于拉模式的。//拉取消息方法 public ConsumerRecordsK, V poll(final Duration timeout)Kafka消费消息是一个不断循环拉取的过程也就是重复的调用poll方法。poll方法是所订阅主题分区上的一组消息。3.2.5 位移提交消费者中的offest来表示消费到分区中某个消息所在的位置。需要持久化保存不然重启后无法知道消费到哪个位置。//获取消费位置(position) public long position(TopicPartition partition) //获取已经提交过的消费位移committed Offset public OffsetAndMetadata committed(TopicPartition partition)position committed offest lastConsumedOffset 1位移提交时机消费者消费获取一批消息如果消费一部分然后异常导致没有位移提交就会导致重复消费。位移提交方式自动提交默认enable.auto.commit配置为true然后定期提交。auto.commit.interval.ms配置周期默认5s.缺点重复消费、消息丢失问题 优点编码简单手动提交enable.auto.commit配置为false。同步提交commitsync 可以按分区提交异步提交commitAsync 可以增加回调方法public void commitSync() public void commitSync(Duration timeout) public void commitSync(final MapTopicPartition, OffsetAndMetadata offsets) public void commitSync(final MapTopicPartition, OffsetAndMetadata offsets, final Duration timeout) public void commitAsync() public void commitAsync(OffsetCommitCallback callback) public void commitAsync(final MapTopicPartition, OffsetAndMetadata offsets, OffsetCommitCallback callback)异步提交回调函数失败如果重试会有先后问题导致重复消费。3.2.6 控制或关闭消费Kafka提供了对消费速度进行控制的方法。通过pause()和resume()方法来分别实现暂停和恢复。3.2.7 指定位移消费新消费者加入时没有可以查找的消费位移。配置auto.offset.restart可以在找不到消费位移时决定从何处开始消费latest 从分区末尾开始也就是下一条earliest:从0开始none:找不到时抛出异常。以上只是找不到时的处理。seek可以指定位移public void seek(TopicPartition partition, long offset)seek只能重置分区的消费位置而拉取哪个分区的消息是poll中实现的所以seek之前必须要先poll通过seek可以跳过或者回溯消息。3.2.8 再均衡分区的所有权从一个消费者转移到另一个消费者的行为。优点高可用伸缩性。可以安全的删除消费组内的消费者或者添加新的消费者。缺点再均衡过程中消费组不可用。消费者状态丢失。 比如消费者还没提交消费位移的时候发生再均衡会导致重复消费。public void subscribe(CollectionString topics) public void subscribe(CollectionString topics, ConsumerRebalanceListener listener) public void subscribe(Pattern pattern) public void subscribe(Pattern pattern, ConsumerRebalanceListener listener)ConsumerRebalanceListener:再均衡监听器用来设定再均衡动作前后的一些准备和收尾动作。public interface ConsumerRebalanceListener { //会在再均衡开始之前和消费者停止读取消息之后被调用 partitions重分配前 void onPartitionsRevoked(CollectionTopicPartition partitions); //在重新分配分区之后和消费者开始读取消息之前被调用。partitions重分配后 void onPartitionsAssigned(CollectionTopicPartition partitions);3.2.9 消费者拦截器Kafka会在poll方法返回结果之前调用拦截器的onConsume方法对消息进行定制化的操作。public interface ConsumerInterceptorK, V extends Configurable, AutoCloseable { //在poll方法返回结果前 public ConsumerRecordsK, V onConsume(ConsumerRecordsK, V records); //在提交完消费位移之后。 public void onCommit(MapTopicPartition, OffsetAndMetadata offsets); public void close();3.2.10 多线程实现生产者是现成安全的但是消费者不是。//KafkaConsumer中 //通过这个方法判断是不是只有一个线程在操作。 相当与一个锁将refcount计数1 private void acquire() { long threadId Thread.currentThread().getId(); if (threadId ! currentThread.get() !currentThread.compareAndSet(NO_CURRENT_THREAD, threadId)) throw new ConcurrentModificationException(KafkaConsumer is not safe for multi-threaded access); refcount.incrementAndGet(); } //释放锁 private void release() { if (refcount.decrementAndGet() 0) currentThread.set(NO_CURRENT_THREAD); }实现方式使用滑动窗口一个方格代表一个批次的消息一个滑动窗口包含若干方法startOffset滑动窗口开始位置endOffset结束位置每当startOffset中的消息被消费完成就能提交这部分位移窗口向前滑动一步。一个方格代表一个线程如果startOffset无法被消费完成悬停一定时间后就可以重试重试失败就转入重试队列再不行就进入死信队列。3.2.11 重要消费参数fetch.min.bytes: 拉取请求中能从Kafka中拉取的最小数据量。如果小于该值会进行等待。fetch.min.bytes: 拉取的最大数据量。fetch.max.wait.ms与fetch.min.bytes参数相关防止一直等待。4 主题与分区4.1 优先副本的选举优先副本AR集合的第一个副本[1,2,0] 优先副本为1。4.2 文件目录一个主题有很多分区分区有很多副本一个副本对应一个目录目录下主要有三类文件。 *.index *.log *.timeindex4.3 日志索引Kafka索引文件以稀疏索引的方式构造消息的索引。每当写入一定量的消息时偏移量索引文件和时间戳索引文件分别增加一个偏移量索引项和时间戳索引项。在索引中使用二分查找法。4.4 日志快速读取日志删除日志压缩相同key的value只保留最新磁盘存储文件只允许追加不允许修改。其实是顺序写磁盘的一种加快了速度。页缓存Kafka不使用Java虚拟机缓存数据使用页缓存。Jvm gc会变慢。零拷贝应用程序直接请求内核磁盘中的数据传输给socket.5、队列就是日志存储按顺序存储然后按位移消费。6、消息传输保障保障层级有三个最多一次至少一次正好一次。6.1如何保证正好一次消息生产者、ack确认机制设置为-1表示消息需要同步ISR才会确认成功。2、重试次数需要大于0。3、开启幂等开启后生产者实力初始化都会分配一个pidpid分区号为key维护一个序列号每次发送消息都会加1broker收到回校验序列号受到过的就不会在接收了。保证了单分区的幂等为了保证多个分区还可以开启事务。同一个事务中事务id相同为了保证事务id不会撞引入epoch事务协调者只会选择epoch最新的那个生产者的事务id旧的抛出异常。实际业务中开启事务回导致性能下降一般消费端自己实现幂等。