23123123
1、入门
1.1 简介
-
什么是事件流
-
事件流是人体中枢神经系统的数字等价物。它是“永远在线”世界的技术基础,在这个世界中,企业越来越多地由软件定义和自动化,软件的用户更多是软件。
-
从技术上讲,事件流是从数据库、传感器、移动设备、云服务和软件应用程序等事件源以事件流的形式实时捕获数据的做法;持久地存储这些事件流以供以后检索;实时和回顾性地操作、处理和响应事件流;并根据需要将事件流路由到不同的目标技术。因此,事件流可确保数据的连续流动和解释,以便正确的信息在正确的时间出现在正确的位置。
-
-
可以使用事件流做什么
- 实时处理支付和金融交易,例如在证券交易所、银行和保险中。
- 实时跟踪和监控汽车、卡车、车队和货物,例如物流和汽车行业。
- 持续捕获和分析来自 IoT 设备或其他设备(例如工厂和风电场)的传感器数据。
- 收集客户互动和订单并立即做出反应,例如在零售、酒店和旅游行业以及移动应用程序中。
- 监测住院病人并预测病情变化,以确保在紧急情况下得到及时治疗。
- 连接、存储和提供公司不同部门产生的数据。
- 作为数据平台、事件驱动架构和微服务的基础。
-
kafka的应用场景
- 要发布(写)和订阅(读)流事件,包括来自其他系统的数据的持续导入/导出的。
- 为了存储持久和可靠的事件流,只要你想要的。
- 在事件发生时或追溯性地处理事件流。
- 所有这些功能都是以分布式、高度可扩展、弹性、容错和安全的方式提供的。Kafka 可以部署在裸机硬件、虚拟机和容器上,也可以部署在本地和云端。
-
kafka如何工作
- 服务器:kakfa作为一个或多个服务器的集群运行,其中一些服务器形成存储层,作为代理。其他服务器将数据作为事件流持续导入和导出,从而将kfka与现有的系统(例如关系型数据库和其他的kafka集群)集成。kafka集群具有高扩展和高容错,其中任何一个服务器出现故障,其他服务器将会接管他的工作,确保持续运行并且不会丢失数据。
- 客户端:编写分布式应用程序和微服务,即使在网络问题和机器故障的情况下,它们也可以并行、大规模和容错方式读取、写入和处理事件流。
-
kafka架构
broker:kafka集群包括一个或多个服务器,这种服务器被称为broker。broker不维护数据的消费状态,直接使用磁盘存储,线性读写,速度快;避免了数据在JVM内存和系统内存之间的复制,减少耗性能的创建对象和垃圾回收。
producer:生产者,是那些向kafka发布(写入)事件的客户端应用程序。
consumer:是那些向kafka订阅(读取和处理)这些事件的客户端应用程序。在 Kafka 中,生产者和消费者之间是完全解耦和不可知的,这是实现 Kafka 众所周知的高可扩展性的关键设计元素。例如,生产者永远不需要等待消费者。
consumer group:消费者组,每个消费者属于某一个消费者组,不分配则为默认组。
Topic: 主题,事件被组织并持久地存储在主题中。主题类似于文件系统中的文件夹,事件就是该文件夹中的文件。一个示例主题名称可以是“支付”。Kafka 中的主题总是多生产者和多订阅者:一个主题可以有零个、一个或多个向其写入事件的生产者,以及零个、一个或多个订阅这些事件的消费者。事件在消费后不会被删除,相反,您可以通过每个主题的配置设置来定义 Kafka 应该保留您的事件多长时间。Kafka 的性能在数据大小方面实际上是恒定的,因此长时间存储数据是完全没问题的。为了使您的数据具有容错性和高可用性,每个主题都可以复制,甚至可以跨地理区域或数据中心复制,以便始终有多个代理拥有数据副本。常见的生产设置是复制因子为 3,即,您的数据将始终存在三个副本。此复制在主题分区级别执行。
Partition:分区,物理上的概念,一个topic具有多个分区。topic逻辑上可以理解为一个queue,每条消息都必须指定它的topic,也就是放在哪个queue里。为了提高kafka的吞吐率,物理上将topic分成一个或多个partition,每个partition在物理上对应一个文件夹,该文件夹下,存储这个partition的所有消息和索引文件。
1.2 用例
- 消息传递
- 网站活动追踪
- 日志聚合、提交日志
- 指标
- 流处理
- 事件溯源
1.3 kafka快速入门
docker pull zookeepr
docker pull wurstmeister/kafka
docker run -d --name zookeeper -p 2181:2181 -t wurstmeister/zookeeper
docker exec -it zookeeper /bin/sh
docker run -d --name kafka -p 9092:9092 -e KAFKA_BROKER_ID=0 -e KAFKA_ZOOKEEPER_CONNECT=175.24.188.175:2181 -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://175.24.188.175:9092 -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 -t wurstmeister/kafka
docker -exec -it kafka /bin/bash
<!--创建主题-->
./kafka-console-producer.sh --broker-list localhost:9092 --topic {topicName}
<!--写一个消息-->
{"datas":[{"name":"jianshu","value":"10"}],"ver":"1.0"}
<!--消费-->
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic {topicName} --from-beginning
<!--创建主题-->
bin/kafka-topics.sh --create --topic {topicname} --bootstrap-server localhost:9092
<!--显示topic使用信息-->
bin/kafka-topics.sh --describe --topic {topicname} --bootstrap-server localhost:9092
2、基础配置参数
-
group.id:
consumer group是kafka提供的可扩展且具有容错性的消费者机制。既然是一个组,那么组内必然可以有多个消费者或消费者实例(consumer instance),它们共享一个公共的ID,即group ID。组内的所有消费者协调在一起来消费订阅主题(subscribed topics)的所有分区(partition)。当然,每个分区只能由同一个消费组内的一个consumer来消费。
-
enable.auto.commit
消费者消费位移的提交方式,true是自动提交,consumer pull消息后自动提交上次之前poll的所有消息位移,若为
false则需要手动提交,即consumer poll出的消息需要手动提交消息位移,提交消息位移的方式有同步提交和异步提交。 -
auto.commit.interval.ms
在enable.auto.commit 为true的情况下,自动提交消费位移的间隔,默认值5000ms。那么消费者会在poll方法调用后每隔5000ms(由auto.commit.interval.ms指定)提交一次位移。和很多其他操作一样,自动提交消费位移也是由poll()方法来驱动的,在调用poll()时,消费者判断是否到达提交时间auto.commit.interval.ms指定的值)如果是则提交上一次poll返回的最大位移。
-
auto.offset.reset
这个参数是针对新的groupid中的消费者而言的,当有新groupid的消费者来消费指定的topic时,对于该参数的配置,会有不同的语义。
auto.offset.reset=latest情况下,新的消费者将会从其他消费者最后消费的offset处开始消费Topic下的消息。
auto.offset.reset= earliest情况下,新的消费者会从该topic最早的消息开始消费。
auto.offset.reset=none情况下,新的消费者加入以后,由于之前不存在offset,则会直接抛出异常。
-
max.poll.records
consumer是通过轮训的方式使用poll()方法不断获取消息的,max.poll.records参数可以限制每次调用poll返回的消息数,默认是500条。
-
max.poll.interval.ms
默认值5分钟,表示若5分钟之内consumer没有消费完上一次poll的消息,也就是在5分钟之内没有调用下次的poll()函数,那么kafka会认为consumer已经宕机,所以会将该consumer踢出consumergroup,紧接着就会发生rebalance,发生rebalance可能会发生重复消费的情况。
需要保证poll出的所有消息消费时间总和不能大于
max.poll.interval.ms,如果大于则会将consumer踢出consumer group,会进行rebalance操作了,每次poll消息的数量不能太大,避免发生rebalance。
3、Topic&Partition
3.1 Topic
-
在kafka中,topic是一个存储消息的逻辑概念,可以认为是一个消息集合。每条消息发送到kafka集群的消息都有一个类别。物理上来说,不同的topic的消息是分开存储的,每个topic可以有多个生产者向它发送消息,也可以有多个消费者去消费其中的消息。
3.2 Partition
- 每个topic可以划分多个分区(每个Topic至少有一个分区),同一topic下的不同分区包含的消息是不同的。第一分区存储可以存储更多的消息,其次是为了提高吞吐量,如果只有一个partition,则所有消息只能存储在该partition内,消费时不管有多少个消费者也只能顺序读取该partition内的消息,如果是多个partition,那么消费者就可以同时从多个partition内并发读取消息。
- 每个消息在被添加到分区时,都会被分配一个offset(称之为偏移量),它是消息在此分区中的唯一编号,kafka通过offset保证消息在分区内的顺序,offset的顺序不跨分区,即kafka只保证在同一个分区内的消息是有序的。
- 在多partition和多consumer的情况下,生产的消息是具有顺序性的,且根据partition的分发策略依次插入到相应的partition中,但是由于kafak只保证同一个partition内的消息输出有序性,所以多partition依次输出的消息顺序并不能保证和生产消息写入的顺序是一样的。
3.3、副本机制
- 我们已经知道Kafka的每个topic都可以分为多个Partition,并且同一topic的多个partition会均匀分布在集群的各个节点下。虽然这种方式能够有效的对数据进行分片,但是对于每个partition来说,都是单点的,当其中一个partition不可用的时候,那么这部分消息就没办法消费。所以kafka为了提高partition的可靠性而提供了副本的概念(Replica),通过副本机制来实现冗余备份。
- 每个分区可以有多个副本,并且在副本集合中会存在一个leader的副本,所有的读写请求都是由leader副本来进行处理。剩余的其他副本都做为follower副本,follower副本会从leader副本同步消息日志。
- 一般情况下,同一个分区的多个副本会被均匀分配到集群中的不同broker上,当leader副本所在的broker出现故障后,可以重新选举新的leader副本继续对外提供服务。通过这样的副本机制来提高kafka集群的可用性。
3.3.1、leader选举机制
三种副本
- leader副本:主副本,每个分区都有一个leader副本,为了保证数据一致性,所有消费者和生产者的请求都会经过该副本处理。
- follower副本:除了leader副本外的其他所有副本都是follower副本,follower副本不处理来自客户端的任何请求,只负责从leader副本同步数据,保证与leader保持一致。如果leader副本发生崩溃,就会从这其中选举出一个leader。
- ISR副本:zookeeper中为每一个partition动态的维护了一个ISR,这个ISR的所有replica都跟上了leader,只有ISR里的成员才有成功leader的可能,ISR副本保存了leader副本以及所有和leader保持同步的follower副本。
副本协同
- 写请求首先由leader副本处理,之后follower副本会从leader上拉取写入的消息,这个过程会有一定的延迟,导致follower副本中保存的消息略少于leader副本,但是只要没有超出阈值都可以容忍。但是如果一个follower副本出现异常,比如宕机、网络断开等原因长时间没有同步到消息,那这个时候,leader就会把它从ISR踢出去。
ISR
-
ISR集合中的副本必须满足两个条件:
- 副本所在节点必须维持着与zookeeper的连接
- 副本最后一条消息的offset与leader副本的最后一条消息的offset之间的差值不能超过指定的阈值。
replica.lag.time.max.ms如果该follower在此时间内一直没有追赶上leader的所有消息,则该follower会被踢出ISR队列。
- follower副本把leader副本的消息全部同步完成,这个时候会更新这个副本的lastCaughtUpTimeMs标识,kafka副本管理器会启动一个副本过期检查的定时任务,这个任务会定期检查当前时间与lastCaughtUpTimeMs的差值是否大于replica.lag.time.max.ms,如果大于,则踢出ISR队列
可用性与一致性
-
在ISR队列中至少有一个follower时,kafka可以确保已经commit的数据不丢失,但如果某个partition的所有replica都宕机了,就无法保证数据不丢失了,这种情况有两种解决方案:
- 等待ISR队列的任何一个replica活过来,并且选为leader,一直等待,不可用时间较长,如果无法活过来,分区将不可用
- 选择第一个活过来的replica(不一定是ISR队列的)作为leader,并不保证已经同步了commit的消息
-
默认情况下Kafka采用第二种策略,即
unclean.leader.election.enable=true,也可设置为false启用第一张策略。
副本数据同步
-
producer在发布消息到某个partition时:
- 通过zookeeper找到该分区的leader副本,producer只将消息发送给leader
- leader会将消息写入其本地log,,每个follower都从leader pull数据。这种方式下,follower与leader存储的数据顺序一致。
- follower收到消息并写入log,向leader发送ACK
- 一旦leader收到了ISR队列中的所有的replica的ACK,该消息就认为已经commit了,leader将增加HW(highWatermark)并且向producer发送ACK
-
LEO:日志末端位移(log end offset),记录了该副本底层日志(log)中下一条消息的位移值。注意是下一条消息!也就是说,如果LEO=10,那么表示该副本保存了10条消息,位移值范围是[0, 9]。
-
HW:即所有follower副本中相对于leader副本最小的LEO值。HW是相对leader副本而言的,其HW值不会大于LEO值。小于等于HW值的所有消息都被认为是“已备份”的(replicated)。
数据可靠性与持久性
-
producer数据不丢失
当向leader发送数据时,可以通过
request.required.acks参数来设置数据可靠性的级别:-
ack=0
producer写入的一条消息会立即返回ack确认消息,不管leader是否同步完成或者ISR中的follower是否同步完成,数据丢失风险大
-
ack=1(默认配置)
producer写入的一条消息会等到leader副本同步完成(不需要等待ISR中的follower副本同步完成)后立即返回cak确认消息。该配置的风险是leader节点宕机了,从ISR中选举一个follower为leader,就有数据丢失的风险。
-
ack=-1
producer写入的一条消息需要等到分区的leader副本和ISR中所有follower副本全部同步完成才会返回ack确认消息,可用性降低。但是也不能保证数据不会丢失,当分区只有一个leader副本时,leader宕机后,数据也会丢失,为了避免只有一个leader副本情况的发生,可以使用
min.insync.replicas来约束,如果ISR中的副本数不够min.insync.replicas设定的值,则会抛出异常。如果由于网络原因producer push失败,也可以设置
retries参数进行重试三种方式保证producer数据不丢失:
设置request.required.acks=1
设置min.insync.replicas参数
设置retries参数
-
-
broker数据不丢失
设置
unclean.leader.election.enable=false,即ISR副本全部宕机后,一直等待ISR中的一个副本活过来,选为leader -
consumer数据不丢失
enable.auto.commit该参数默认为true,表明consumer在下次poll消息时自动提交上次poll出的所有消息的消费位移,如果设置为false,则需要用户手动提交手动提交所有消息的消费位移。
重复消费、消息丢失
-
enable.auto.commit设置为true时,会有消息重复消费和消息丢失的场景。- 当应用端消费消息时,还没有提交消费位移的时候,此时kafka出现宕机,那么在kafka恢复之后,这些消息将会重新被消费一遍,这就造成了重复消费。
- consumer第一次poll出n条消息进行消费,达到
auto.commit.interval.ms时间后,cosumer会进行下一次poll并提交上次poll出的n条消息的消费位移。如果第一次poll出的n条消息客户端还没有消费完,此时客户端宕机了,当客户端重启后,,将会从第二次poll的位置开始拉取消息,从而丢失第一次未提交消费位移的消息,这就造成了数据丢失。
-
enable.auto.commit设置为false时,可以避免消息丢失,而不能保证消息重复消费
保证exactly once
https://blog.csdn.net/w1992wishes/article/details/89502956
3.4 确定分区数
- 创建一个只有一个分区的topic,然后测试这个topic的producer吞吐量Tp和consumer吞吐量Tc,单位可以使MB/s,假设总的吞吐量是Tt,那么分区数=Tt/max(Tp,Tc)
- 测试producer的吞吐量,直接发送消息就可以了。consumer的吞吐量Tc通常与应用的关系更大, 因为Tc的值取决于你拿到消息之后执行什么操作,因此Tc的测试通常也要麻烦一些。
3.5 Topic&Partition的存储
Partition是以文件的形式存储在文件系统中,比如创建一个名为firstTopic的topic,其中有3个partition,那么在kafka的数据目录(/tmp/kafka-log)中就有3个目录,命名规则是<topic_name>-<partition_id>。
4、消息
4.1、消息格式
- 一个Kafka的Message由一个固定长度的header和一个变长的消息体body组成
- header部分由一个字节的magic(文件格式)和四个字节的CRC32(用于判断body消息体是否正常,是否丢包,数据不一样CRC32算出来的数字也是不一样的)构成。
- 当magic的值为1的时候,会在magic和crc32之间多一个字节的数据:attributes(保存一些相关属性,比如是否压缩、压缩格式等等);如果magic的值为0,那么不存在attributes属性
- body是由N个字节构成的一个消息体,包含了具体的key/value消息
4.2、Log消息格式
-
存储在磁盘的日志采用不同于Producer发送的消息格式
-
每个日志文件都是一个“log entries”序列
- 每一个log entry包含一个四字节整型数(message长度,值为1+4+N)
- 一个字节的magic
- 四个字节的CRC32值
- 最终是N个字节的消息数据。每条消息都有一个当前Partition下唯一的64位offset
-
其实这个log entries也不是一个文件,是一个index(索引文件)和一个log日志文件
4.3、生产者消息分发策略
4.3.1、生产者push消息的模式
- Kafka的发送模式由producer端的配置参数
producer.type来设置,这个参数指定了在后台线程中消息的发送方式是同步的还是异步的,默认是同步的,即producer.type=sync - 如果设置成异步的方式,即
producer.type=async,producer可以以batch的形式push数据,就是将消息按批量的方式发送,这样会极大的提高broker的性能,但是这样会增加丢失数据的风险。
4.3.2、高可靠性配置
- 分区副本, 你可以创建分区副本来提升数据的可靠性,避免数据丢失,但是分区数过多也会带来性能上的开销,一般来说,3个副本就能满足对大部分场景的可靠性要求
- broker的配置:leader的选举条件
unclean.leader.election.enable=false - producer配置:
request.required.acks=-1,producer.type=sync
4.3.3、消息分发策略
-
消息是kafka中最基础的数据单元。在发送一条消息时,我们可以指定这个key,那么producer会根据key和partition机制来判断当前这条消息应该发送并存储到哪个partition中。我们可以根据需要进行扩展producer的partition机制。
-
生产者在将消息发送到某个Topic ,需要经过拦截器、序列化器和分区器(
Partitioner)的一系列作用之后才能发送到对应的Broker,在发往Broker之前是需要确定它所发往的分区。- 如果消息
ProducerRecord指定了partition字段,那么就不需要分区器 - 如果消息
ProducerRecord没有指定partition字段,那么就需要依赖分区器,根据
- 如果消息
public class ProducerRecord<K, V> {
// 该消息需要发往的主题
private final String topic;
// 该消息需要发往的主题中的某个分区,如果该字段有值,则分区器不起作用,直接发往指定的分区
// 如果该值为null,则利用分区器进行分区的选择
private final Integer partition;
private final Headers headers;
// 如果partition字段为null,则使用分区器进行分区选择时会用到该key字段,该值可为空
private final K key;
private final V value;
private final Long timestamp;
- Kafka 中提供的默认分区器是
DefaultPartitioner,它实现了Partitioner接口(用户可以实现这个接口来自定义分区器),其中的partition方法就是用来实现具体的分区分配逻辑:- 如果在发消息的时候指定了分区,则消息投递到指定的分区。
- 如果没有指定分区,但是消息的key不为空,则使用称之为
murmur的Hash算法来计算分区分配。 - 如果既没有指定分区,且消息的key也是空,则用轮询的方式选择一个分区。
public class DefaultPartitioner implements Partitioner {
private final ConcurrentMap<String, AtomicInteger> topicCounterMap = new ConcurrentHashMap<>();
public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
// 首先通过cluster从元数据中获取topic所有的分区信息
List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
// 拿到该topic的分区数
int numPartitions = partitions.size();
// 如果消息记录中没有指定key
if (keyBytes == null) {
// 则获取一个自增的值
int nextValue = nextValue(topic);
// 通过cluster拿到所有可用的分区(可用的分区这里指的是该分区存在首领副本)
List<PartitionInfo> availablePartitions = cluster.availablePartitionsForTopic(topic);
// 如果该topic存在可用的分区
if (availablePartitions.size() > 0) {
// 那么将nextValue转成正数之后对可用分区数进行取余
int part = Utils.toPositive(nextValue) % availablePartitions.size();
// 然后从可用分区中返回一个分区
return availablePartitions.get(part).partition();
} else { // 如果不存在可用的分区
// 那么就从所有不可用的分区中通过取余的方式返回一个不可用的分区
return Utils.toPositive(nextValue) % numPartitions;
}
} else { // 如果消息记录中指定了key
// 则使用该key进行hash操作,然后对所有的分区数进行取余操作,这里的hash算法采用的是murmur2算法,然后再转成正数
//toPositive方法很简单,直接将给定的参数与0X7FFFFFFF进行逻辑与操作。
return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
}
}
// nextValue方法可以理解为是在消息记录中没有指定key的情况下,需要生成一个数用来代替key的hash值
// 方法就是最开始先生成一个随机数,之后在这个随机数的基础上每次请求时均进行+1的操作
private int nextValue(String topic) {
// 每个topic都对应着一个计数
AtomicInteger counter = topicCounterMap.get(topic);
if (null == counter) { // 如果是第一次,该topic还没有对应的计数
// 那么先生成一个随机数
counter = new AtomicInteger(ThreadLocalRandom.current().nextInt());
// 然后将该随机数与topic对应起来存入map中
AtomicInteger currentCounter = topicCounterMap.putIfAbsent(topic, counter);
if (currentCounter != null) {
// 之后把这个随机数返回
counter = currentCounter;
}
}
// 一旦存入了随机数之后,后续的请求均在该随机数的基础上+1之后进行返回
return counter.getAndIncrement();
}
4.4、Offset
4.4.1、什么是offset
- 每个topic可以划分多个分区(每个Topic至少有一个分区),同一topic下的不同分区包含的消息是不同的。每个消息在被添加到分区时,都会被分配一个offset(称之为偏移量),它是消息在此分区中的唯一编号,kafka通过offset保证消息在分区内的顺序,offset的顺序不跨分区,即kafka只保证在同一个分区内的消息是有序的。对于应用层的消费来说,每次消费一个消息并且提交以后,会保存当前消费到的最近的一个offset。
4.4.2、offset在哪里
- 在kafka中,提供了一个consumer_offsets_* 的一个topic,把offset信息写入到这个topic中。它保存了每个consumer group某一时刻提交的offset信息,默认有50个分区。
4.5、消费者消息分配策略
-
消费者以组的名义订阅主题,主题有多个分区,消费者组中有多个消费者实例,同一时刻,一条消息只能被组中的一个消费者实例消费。
- 如果分区数大于或者等于组中的消费者实例数,一个消费者会负责多个分区。
- 如果分区数小于组中的消费者实例数,有些消费者将处于空闲状态并且无法接收消息。
-
如果多个消费者负责同一个分区,那么就意味着两个消费者同时读取分区的消息,由于消费者自己可以控制读取消息的Offset,就有可能C1才读到2,而C1读到1,C1还没处理完,C2已经读到3了,这就相当于多线程读取同一个消息,会造成消息处理的重复,且不能保证消息的顺序。
4.4.1、RangeAssignor
-
Range:默认分配策略
- 首先,将分区按数字顺序排序,消费者按名称的字典顺序排序。
- 然后,用分区总数除以消费者总数。如果能够除尽,平均分配;若除不尽,则位于排序前面的消费者将多负责一个分区。
-
假设,有1个主题、10个分区、3个消费者线程, 10 / 3 = 3,而且除不尽,那么消费者C1将会多消费一个分区,分配结果是:
- C1 将消费T1主题的0、1、2、3分区。
- C2 将消费T1主题的4、5、6分区。
- C3 将消费T1主题的7、8、9分区
-
假设,有11个分区,分配结果是:
- C1 将消费T1主题的0、1、2、3分区。
- C2 将消费T1主题的4、5、 6、7分区。
- C3 将消费T1主题的8、9、10分区。
-
假如我们有2个主题(T1和T2),分别有10个分区
- C1 将消费 T1主题的 0, 1, 2, 3 分区以及 T2主题的 0, 1, 2, 3分区
- C2 将消费 T1主题的 4, 5, 6 分区以及 T2主题的 4, 5, 6分区
- C3 将消费 T1主题的 7, 8, 9 分区以及 T2主题的 7, 8, 9分区
public Map<String, List<TopicPartition>> assign(Map<String, Integer> partitionsPerTopic,
Map<String, Subscription> subscriptions) {
// 主题与消费者的映射
Map<String, List<String>> consumersPerTopic = consumersPerTopic(subscriptions);
Map<String, List<TopicPartition>> assignment = new HashMap<>();
for (String memberId : subscriptions.keySet())
assignment.put(memberId, new ArrayList<TopicPartition>());
for (Map.Entry<String, List<String>> topicEntry : consumersPerTopic.entrySet()) {
String topic = topicEntry.getKey(); // 主题
List<String> consumersForTopic = topicEntry.getValue(); // 消费者列表
// partitionsPerTopic表示主题和分区数的映射
// 获取主题下有多少个分区
Integer numPartitionsForTopic = partitionsPerTopic.get(topic);
if (numPartitionsForTopic == null)
continue;
// 消费者按字典序排序
Collections.sort(consumersForTopic);
// 分区数量除以消费者数量
int numPartitionsPerConsumer = numPartitionsForTopic / consumersForTopic.size();
// 取模,余数就是额外的分区
int consumersWithExtraPartition = numPartitionsForTopic % consumersForTopic.size();
List<TopicPartition> partitions = AbstractPartitionAssignor.partitions(topic, numPartitionsForTopic);
for (int i = 0, n = consumersForTopic.size(); i < n; i++) {
int start = numPartitionsPerConsumer * i + Math.min(i, consumersWithExtraPartition);
int length = numPartitionsPerConsumer + (i + 1 > consumersWithExtraPartition ? 0 : 1);
// 分配分区
assignment.get(consumersForTopic.get(i)).addAll(partitions.subList(start, start + length));
}
}
return assignment;
}
4.4.2、RoundRobinAssignor
-
RoundRobin基于轮询算法
-
轮询分区策略是把所有partition和所有consumer线程都列出来,然后按照hashcode进行排序。注意上一种range分区是针对每一个topic而言的,而轮训分区是相对于所有的partition和consumer而言的,最后通过轮询算法分配partition给消费线程。如果消费组内,所有消费者订阅的Topic列表是相同的(每个消费者都订阅了相同的Topic),那么分配结果是尽量均衡的(消费者之间分配到的分区数的差值不会超过1)。如果订阅的Topic列表是不同的,那么分配结果是不保证“尽量均衡”的,因为某些消费者不参与一些Topic的分配。
-
举例:假如按照 hashCode 排序完的topicpartitions组依次为 T1-5, T1-3, T1-0, T1-8, T1-2, T1-1, T1-4, T1-7, T1-6, T1-9,我们的消费者线程排序为 C1-0, C1-1, C2-0, C2-1,最后分区分配的结果为:
- C1-0 将消费 T1-5, T1-2, T1-6 分区
- C1-1 将消费 T1-3, T1-1, T1-9 分区
- C2-0 将消费 T1-0, T1-4 分区
- C2-1 将消费 T1-8, T1-7 分区
-
使用轮询分区策略必须满足两个条件
- 每个主题的消费者实例具有相同数量的流
- 每个消费者订阅的主题必须是相同的
-
相对于RangeAssignor,在订阅多个Topic的情况下,RoundRobinAssignor的方式能消费者之间尽量均衡的分配到分区(分配到的分区数的差值不会超过1,RangeAssignor的分配策略可能随着订阅的Topic越来越多,差值越来越大) -
对于订阅组内消费者订阅Topic不一致的情况:假设有三个消费者分别为C1、C2、C3,有3个Topic T1、T2、T3,分别拥有1、2、3个分区,并且C1订阅T1,C2订阅T1和T2,C3订阅T1、T2、T3,那么RoundRobinAssignor的分配结果如下:

4.4.3、StickyAssignor
-
背景
- 尽管RoundRobinAssignor已经在RangeAssignor上做了一些优化来更均衡的分配分区,但是在一些情况下依旧会产生严重的分配偏差,比如消费组中订阅的Topic列表不相同的情况下(这个情况可能更多的发生在发布阶段,但是这真的是一个问题吗?——可以参照Kafka官方的说明:KIP-49 Fair Partition Assignment Strategy)。更核心的问题是无论是RangeAssignor,还是RoundRobinAssignor,当前的分区分配算法都没有考虑上一次的分配结果。显然,在执行一次新的分配之前,如果能考虑到上一次分配的结果,尽量少的调整分区分配的变动,显然是能节省很多开销的。
-
粘性策略:从字面意义上看,Sticky是“粘性的”,可以理解为分配结果是带“粘性的”——每一次分配变更相对上一次分配做最少的变动(上一次的结果是有粘性的),其目标有两点:
- 分区的分配尽可能的均匀
- 分区的分配尽可能和上次分配保持相同,也就是
rebalance之后分区的分配尽量和之前的分区分配相同。 - 当这两个目标发生冲突时,优先保证第一个目标。第一个目标是每个分配算法都尽量尝试去完成的,而第二个目标才真正体现出StickyAssignor特性的。
-
举例1:
- 有3个Consumer:C0、C1、C2
- 有4个Topic:T0、T1、T2、T3,每个Topic有2个分区
- 所有Consumer都订阅了这4个分区
-
举例2:
- 有3个Consumer:C0、C1、C2
- 3个Topic:T0、T1、T2,它们分别有1、2、3个分区
- C0订阅T0;C1订阅T0、T1;C2订阅T0、T1、T2
4.6、rebalance
4.5.1、触发场景
-
当出现以下几种情况时,kafka 会进行一次分区分配操作,也就是 kafka consumer的 rebalance
- consumer增加或删除会触发consumer group的rebalance
- 订阅的 Topic 个数发生变化。
- 订阅 Topic 的分区数发生变化。
-
rebalance 发生时,group 下所有 consumer 实例都会协调在一起共同参与,kafka 能够保证尽量达到最公平的分配。但是 rebalance 过程对 consumer group 会造成比较严重的影响。在 rebalance 的过程中 consumer group 下的所有消费者实例都会停止工作,等待 rebalance 过程完成。
4.5.3、metadata
-
kafka集群的metadata包括:
-
所有broker的信息: ip和port;
-
所有topic的信息: topic name, partition数量, 每个partition对应的broker、leader、 isr, replica集合等
-
kafka集群的每一台broker都缓存了整个集群的metadata, 当broker或某一个topic的metadata信息发生变化时, 集群的controller都会感知到作相应的状态转换, 同时把发生变化的新的metadata信息广播到所有的broker
4.5.2、coordinator
-
Group Coordinator是一个服务,每个Broker在启动的时候都会启动一个该服务。Group Coordinator的作用是用来存储Group的相关Meta信息,并将对应Partition的Offset信息记录到Kafka内置Topic(__consumer_offsets)中。
-
Kafka在0.9之前是基于Zookeeper来存储Partition的Offset信息(consumers/{group}/offsets/{topic}/{partition}),因为ZK并不适用于频繁的写操作,所以在0.9之后通过内置Topic的方式来记录对应Partition的Offset。
-
确定coordinator:消费者向kafka集群中的任意一个broker发送一个
GroupCoordinatorRequest请求,服务端会返回一个负载最小的broker节点的id,并将该broker设置为coordinator -
JoinGroup
在rebalance之前,需要保证coordinator是已经确定好了的,整个rebalance的过程分为两个步骤,Join和Sync。
-
Join:表示加入到consumer group中,在这一步中,所有的成员都会向coordinator发送joinGroup的请求。一旦所有成员都发送了joinGroup请求,那么coordinator会选择一个consumer担任leader角色,并把组成员信息和订阅信息发送消费者。
leader选举算法比较简单,如果消费组内没有leader,那么第一个加入消费组的消费者就是消费者leader,如果这个时候leader消费者退出了消费组,那么重新选举一个leader。
-
每个消费者都可以设置自己的分区分配策略,对于消费组而言,会从各个消费者上报过来的分区分配策略中选举一个彼此都赞同的策略来实现整体的分区分配,这个"赞同"的规则是,消费组内的各个消费者会通过投票来决定
-
在joingroup阶段,每个consumer都会把自己支持的分区分配策略发送到coordinator
-
coordinator收集到所有消费者的分配策略组成一个候选集
-
每个消费者需要从候选集里找出一个自己支持的策略,并且为这个策略投票
-
最终计算候选集中各个策略的选票数,票数最多的就是当前消费组的分配策略
-
-
Sync:完成分区分配之后,就进入了Synchronizing Group State阶段,主要逻辑是向GroupCoordinator发送SyncGroupRequest请求,并且处理SyncGroupResponse响应,简单来说,就是leader将消费者对应的partition分配方案同步给consumer group 中的所有consumer
每个消费者都会向coordinator发送syncgroup请求,不过只有leader节点会发送分配方案,其他消费者只是打打酱油而已。当leader把方案发给coordinator以后,coordinator会把结果设置到SyncGroupResponse中。这样所有成员都知道自己应该消费哪个分区。consumer group的分区分配方案是在客户端执行的!Kafka将这个权利下放给客户端主要是因为这样做可以有更好的灵活性
-
4.7、消息持久化
- kafka是使用日志文件的方式来保存生产者和发送者的消息,每条消息都有一个offset值来表示它在分区中的偏移量。Kafka中存储的一般都是海量的消息数据,为了避免日志文件过大,Log并不是直接对应在一个磁盘上的日志文件,而是对应磁盘上的一个目录,这个目录命名规则是<topic_name>_<partition_id>
4.7.1、存储机制
- 一个topic的多个分区在物理磁盘上的保存路径在/tmp/kafka-logs/topic_partition,包含日志文件、索引文件和时间索引文件
- kafka是通过分段的方式将Log分为多个LogSegment,LogSegment是一个逻辑上的概念,一个LogSegment对应磁盘上的一个日志文件和一个索引文件,其中日志文件是用来记录消息的,索引文件是用来保存消息的索引。
LogSegment
- 假设kafka以partition为最小存储单位,那么我们可以想象当kafka producer不断发送消息,必然会引起partition文件的无线扩张,这样对于消息文件的维护以及被消费的消息的清理带来非常大的挑战,所以kafka以segment为单位又把partition进行细分。每个partition相当于一个巨型文件被平均分配到多个大小相等的segment数据文件中每个segment文件中的消息不一定相等),这种特性方便已经被消费的消息的清理,提高磁盘的利用率。
log.segment.bytes可以设置分段大小,默认是1G- segment文件命名规则:partion全局的第一个segment从0开始,后续每个segment文件名为上一个segment文件最后一条消息的offset值进行递增。数值最大为64位long大小,20位数字字符长度,没有数字用0填充
- index中存储了索引以及物理偏移量, log存储了消息的内容。举个简单的案例来说,在索引文件中以[4053,80899]为例,在log文件中,对应的是第4053条记录,物理偏移量(position)为80899
4.7.2、查找message
- 根据offset的值,查找segment段中的index索引文件。由于索引文件命名是以上一个文件的最后一个offset进行命名的,,所以,使用二分查找算法能够根据offset快速定位到指定的索引文件。
- 找到索引文件后,根据offset进行定位,找到索引文件中的符合范围的索引。
- 得到position以后,再到对应的log文件中,从position出发开始查找offset对应的消息,将每条消息的offset与目标offset进行比较,直到找到消息
举例:比如说,我们要查找offset=2490这条消息,那么先找到00000000000000000000.index,然后找到[2487,49111]这个索引,再到log文件中,根据49111这个position开始查找,比较每条消息的offset是否大于等于2490,最后查找到对应的消息以后返回
4.7.3、日志清除、日志压缩
-
日志清除策略
- 根据消息的保留时间,当消息在kafka中保存的时间超过了指定的时间,就会触发清理过程
- 根据topic存储的数据大小,当topic所占的日志文件大小大于一定的阀值,则可以开始删除最旧的消息。kafka会启动一个后台线程,定期检查是否存在可以删除的消息
- 通过
log.retention.bytes和log.retention.hours这两个参数来设置,当其中任意一个达到要求,都会执行删除。默认保留时间:7天
-
日志压缩策略
-
Kafka还提供了“日志压缩(Log Compaction)”功能,通过这个功能可以有效的减少日志文件的大小,缓解磁盘紧张的情况。
-
在很多实际场景中,消息的key和value的值之间的对应关系是不断变化的,就像数据库中的数据会不断被修改一样,消费者只关心key对应的最新的value。
-
因此,我们可以开启kafka的日志压缩功能,服务端会在后台启动启动Cleaner线程池,定期将相同的key进行合并,只保留最新的value值。
-
4.7.4、磁盘存储
-
顺序读写
- 磁盘读写有两种方式:顺序读写或者随机读写。在顺序读写的情况下,磁盘的顺序读写速度和内存持平。
- 因为磁盘是机械结构,每次读写都会寻址->写入,其中寻址是一个“机械动作”。为了提高读写磁盘的速度,Kafka 就是使用顺序 I/O。
- Kafka 利用了一种分段式的、只追加 (Append-Only) 的日志,基本上把自身的读写操作限制为顺序 I/O,也就使得它在各种存储介质上能有很快的速度。
-
零拷贝
-
页缓存
-
Kafka 接收来自 socket buffer 的网络数据,应用进程不需要中间处理、直接进行持久化时,可以使用mmap 内存文件映射。
-
Memory Mapped Files
-
简称
mmap,简单描述其作用就是:将磁盘文件映射到内存,用户通过修改内存就能修改磁盘文件。 -
它的工作原理是直接利用操作系统的 Page 来实现磁盘文件到物理内存的直接映射。完成映射之后你对物理内存的操作会被同步到硬盘上(操作系统在适当的时候)。
-
mmap 也有一个很明显的缺陷:不可靠,写到 mmap 中的数据并没有被真正的写到硬盘,操作系统会在程序主动调用 flush 的时候才把数据真正的写到硬盘。
Kafka 提供了一个参数
producer.type来控制是不是主动 flush:- 如果 Kafka 写入到 mmap 之后就立即 flush,然后再返回 Producer 叫同步(sync);
- 写入 mmap 之后立即返回 Producer 不调用 flush 叫异步(async)。
-
-
当一个进程准备读取磁盘上的文件内容时
- 操作系统会先查看待读取的数据所在的页(page)是否在页缓存(pagecache)中,如果存在(命中) 则直接返回数据,从而避免了对物理磁盘的 I/O 操作;
- 如果没有命中,则操作系统会向磁盘发起读取请求并将读取的数据页存入页缓存,之后再将数据返回给进程。
-
如果一个进程需要将数据写入磁盘
- 操作系统也会检测数据对应的页是否在页缓存中,如果不存在,则会先在页缓存中添加相应的页,最后将数据写入对应的页。
- 被修改过后的页也就变成了脏页,操作系统会在合适的时间把脏页中的数据写入磁盘,以保持数据的一致性。
-
5、一些参数
5.1、生产者
-
batch.size:为了提升吞吐量,生产者传输消息采用分批次的做法,batch.size批次大小,默认16K -
acks:消息验证,0,1,-1 -
retries:重试次数,默认0,发送异常时不进行任何重试动作。- 设置重试可以很好地应对那些瞬时错误,因此推荐用户设置该参数为一个大于 0 的值 。 只不过在考虑
retries 的设置时,有两点需要着重注意 :- 重试可能造成消息的重复发送,为了应对这一风险, Kafka 要求用户在 consumer 端必须执行去重处理 。 令人欣喜的是,社区己于 0.11.0.0 版本开始支持“精确一次”处理语义,从设计上避免了类似的问题 。
- 重试可能造成消息的乱序,当前 producer 会将多个消息发送请求(默认是 5 个) 缓存在内存中,如果由于某种原因发生了消息发送的重试,就可能造成消息流的乱序 。为了避免乱序发生, Java 版本 producer 提供了 max.in.flight.requets.per.connection 参数 ,一旦用户将此参数设置成 1, producer 将确保某一时刻只能发送一个请求
- 设置重试可以很好地应对那些瞬时错误,因此推荐用户设置该参数为一个大于 0 的值 。 只不过在考虑
-
retry.backoff.ms:两次重试之间的时间间隔,避免无效的频繁重试,默认100ms -
linger.ms:生产者会在producerBatch被填满或等待时间超过linger.ms值时发送出去,增大这个参数会增加消息的延迟,但能提升一定的吞吐量 -
buffer.memory:缓冲区大小,默认32M -
compression.type:压缩方式,默认值none,默认消息不会被压缩 -
max.block.ms:该配置控制 KafkaProducer.send() 和 KafkaProducer.partitionsFor() 将阻塞多长时间,默认60S -
max.request.size:生产者能发送的消息的最大值,默认1M -
request.timeout.ms:当 producer 发送请求给 broker 后 , broker 需要在规定 的时 间范围 内 将处理结果返还给producer 。 这段时间便是由该参数控制的,默认是 30 秒 。这就是说,如果 broker 在 30 秒内都没有给 producer 发送响应,那么 producer 就会认为该请求超时了,并在回调函数中显式地抛出TimeoutException 异常交由用户处理 。
默认的 30 秒对于一般的情况而言是足够的 , 但如果 producer 发送的负载很大 , 超时的情况就很容易碰到,此时就应该适当调整该参数值。
-
max.in.flight.requests.per.connection:每个发送数据的网络连接对并未接收到响应的消息的最大数。默认值是5。producer向各个服务器发送数据都会建立不同的网络连接,然后开始发送数据,假如现在我们的max.in.flight.requests.per.connection设置成默认值5,发送了1,2,3,4,5,这服务器都没给我们返回响应,那消息6我们就不能继续再发了。 -
enable.idempotence:幂等性开启,默认为false -
receive.buffer.bytes:socket介绍消息缓冲区的大小,默认32K,如果设置为-1,则使用操作系统的默认值 -
send.buffer.bytes:socket发送消息缓存区的大小,默认128K
5.2、消费者
-
session.timeout.ms:consumer group检查组内成员发送崩溃的时间,假设设置为5分钟,那么当某个group成员突然崩溃了,coordinator可能5分钟之后才会感知到这个崩溃。在实际使用中,可以设为比较小的值让coordinator能够更快的检查consumer的崩溃,避免消息滞后(consumer lag)。目前参数默认值为10S -
max.poll.interval.ms:consumer处理消息逻辑的最大时间,倘若consumer两次poll之间的时间间隔超过该设置时间,coordinator就会认为consumer已经跟不上组内成员的消费进度,coordinator会将该consumer踢出consumer group,该consumer消费的分区会分配给其他consumer,也会触发rebalance假设用户的业务场景中消息处理逻辑是把消息“落地”到远程数据库中,且这个过程平均处理时间是 2 分钟,那么用户仅需要将 max.poll.interval.ms 设置为稍稍大于 2 分钟的值即可
**通过将该参数设置成实际的逻辑处理时间再结合较低的 session.timeout.ms 参数值,consumer group 既实现了快速的 consumer 崩溃检测,也保证了复杂的事件处理逻辑不会造成不必要的 rebalance **
-
auto.offset.reset:指定了无位移信息或位移越界(即consumer要消费的消息位移不在消息日志的合理区间范围)时kafka的应对策略。- earliest:指定从最早的位移开始消费,最早位移不一定为0
- latest:指定从最新处位移开始消费
- none:指定如果未出现无位移信息或者位移越界,抛出异常,使用甚少
-
enable.auto.commit:指定consumer是否主动提交,若设置为true,主动提交,否则用户需要手动提交位移,对于exactly once语义,最好设为false。 -
fetch.max.bytes:他指定了consumer端单次能够获取数据的最大字节数,若实际业务消息很大,则必须要设置该参数为一个很大的值,否则consumer无法消费这些消息 -
max.poll.records:该参数控制单次 poll 调用返回的最大消息数 ,如果用户发现 consumer 端的瓶颈在 poll 速度太慢,可以适当地增加该参数的值。默认500条 -
heartbeat.interval.ms:该值必须小于session.timeout.ms,如果 consumer 在 session.timeout.ms 这段时间内都不发送心跳, coordinator 就会认为它已经 dead,因此也就没有必要让它知晓 coordinator 的决定了 。new KafkaConsumer对象后,在while true循环中执行consumer.poll拉取消息这个过程中,其实背后是有2个线程的,即一个kafka consumer实例包含2个线程:
一个是heartbeat 线程,另一个是processing线程,processing线程可理解为调用consumer.poll方法执行消息处理逻辑的线程,而heartbeat线程是一个后台线程,对程序员是"隐藏不见"的。如果消息处理逻辑很复杂,比如说需要处理5min,那么 max.poll.interval.ms可设置成比5min大一点的值。
heartbeat 线程则和上面提到的参数 heartbeat.interval.ms有关,heartbeat线程每隔heartbeat.interval.ms向coordinator发送一个心跳包,证明自己还活着。只要 heartbeat线程 在 session.timeout.ms 时间内向 coordinator发送过心跳包,那么group coordinator就认为当前的kafka consumer是活着的。
-
connection.max.idle.ms:经常有用户抱怨在生产环境下周期性地观测到请求平均处理时间在飘升,这很有可能是因为 Kafka 会定期地关闭空闲 Socket 连接导致下次 consumer 处理请求时需要重新创建连向 broker 的 Socket 连接 。 当前默认值是 9 分钟,如果用户实际环境中不在乎这些 Socket 资源开销,比较推荐设置该参数值为-1 ,即不要关闭这些空闲连接 。
6、源码分析
6.1、Producer
整体流程图

waitOnMetadata:

first make sure the metadata for the topic is available:确保主题的元数据可用
- 我们在发送消息前就是通过这个waitOnMetadata来同步等待元数据的拉取的。maxBlockTimeMs是指最多等待这个拉取过程多久,因为这个拉取过程进行时代码是阻塞在这里的,所以我们必须设置一个时间限制来放行。
- 然后计算了一下剩余时间,获取更新的元数据
5.1.1、Serializer序列化器
- kafka序列化消息是在生产端,序列化后,消息才能网络传输。
- 属性keySerializer和valueSerializer就是key和value指定的序列化方式,无论是key还是value序列化和反序列化实现都是一样的。
5.1.2、Partitioner分区器
public class ProducerRecord<K, V> {
// 该消息需要发往的主题
private final String topic;
// 该消息需要发往的主题中的某个分区,如果该字段有值,则分区器不起作用,直接发往指定的分区
// 如果该值为null,则利用分区器进行分区的选择
private final Integer partition;
private final Headers headers;
// 如果partition字段为null,则使用分区器进行分区选择时会用到该key字段,该值可为空
private final K key;
private final V value;
private final Long timestamp;
在发往Broker之前是需要确定它所发往的分区。
- 如果消息
ProducerRecord指定了partition字段,那么就不需要分区器 - 如果消息
ProducerRecord没有指定partition字段,那么就需要依赖分区器
Kafka 中提供的默认分区器是 DefaultPartitioner,它实现了Partitioner接口(用户可以实现这个接口来自定义分区器),其中的partition方法就是用来实现具体的分区分配逻辑:
- 如果在发消息的时候指定了分区,则消息投递到指定的分区。
- 如果没有指定分区,但是消息的key不为空,则使用称之为
murmur的Hash算法来计算分区分配。 - 如果既没有指定分区,且消息的key也是空,则用轮询的方式选择一个分区。
5.1.3、RecordAccumulator
-
KafkaProduer可以有同步和异步两种方式发送消息,其实两者的底层实现相同,都是通过异步方式实现的。
-
主线程调用send方法发送消息的时候,先将消息放到RecordAccumulator中暂存,然后主线程就可以返回了,这个时候消息并没有真正的发送给Kafka,而是放在了RecordAccumulator中,
-
之后,主线程通过send方法不断的向RecordAccumulator里面不断追加消息,当达到一定的条件,会唤醒Sender线程发送RecordAccumulator里面的消息。
-
RecordAccumulator中有一个以TopicPartition为key的ConcurrentMap,每个value是ArrayDeque
,其中缓存了发往对应的TopicPartition的消息。每个RecordBatch拥有一个MemoryRecords对象的引用,他才是消息最终存放的地方。 -
send的时候,KafkaProducer把消息放到本地的消息队列RecordAccumulator,然后一个后台线程Sender不断循环,把消息发给Kafka集群。要实现这个,还得有一个前提条件:就是KafkaProducer/Sender都需要获取集群的配置信息Metadata。
-
客户端从metadata中获取topic的partition信息, 知道leader是谁, 才可以发送和消费msg

producer 有两个线程,一个是 main 线程用于把消息放到 RecordAccumulator 寄存器中寄存。另一个线程是 sender线程会通过 IO 和 kafka server 进行交互发送消息
5.1.4、Metadata
- 简单说, kafka集群的metadata包括:
- 所有broker的信息: ip和port;
- 所有topic的信息: topic name, partition数量, 每个partition的leader, isr, replica集合等
- kafka集群的每一台broker都缓存了整个集群的metadata,当broker或某一个topic的metadata信息发生变化时, 集群的controller都会感知到作相应的状态转换, 同时把发生变化的新的metadata信息广播到所有的broker
- 可以从任意broker获取想要的metadata
5.1.5、Sender

-
在 Sender 线程中有三个组件:
-
KafkaClient:它是一个接口,实现它的是 NetworkClient,它用于把消息封装成clientRequest并且放入 selector 中的 send 字段,所以说真正与 sever 进行 io 操作的是 selector,它采用了 java 中的 nio
-
Metadata:用于获取集群中的元数据,更新等操作
-
RecordAccumulator:用于获取预备好的集合,会调用 ready 方法。
-
5.1.6、无消息丢失和消息乱序
-
KafkaProducer.send 方法仅仅把消息放入缓冲区中,由一个专属 IO 线程负责从缓冲区中提取消息井封装进消息 batch 中,然后发送出去。显然,这个过程中存在着数据丢失的窗口:若 IO线程发送之前 producer 崩溃,则存储缓冲区中的消息全部丢失了。
-
producer 的另一个问题就是消息的乱序:
producer.send(record1); producer.send(record2);若此时由于某些原因(比如瞬时的网络抖动〉导致 record I 未发送成功,同时 Kafka 又配置了重试机制以及 max.in.flight.requests.per.connection 大于 1 (默认值是 5 ),那么 producer 重试 record1成功后, record1在日志中的位置反而位于 record2 之后,这样造成了消息的乱序。
-
同步发送:
-
既然异步发送可能丢失数据,改成同步发送似乎是一个不错的主意,但是性能会很差
producer.send(record1).get();
-
-
无消息丢失配置:
-
producer端配置
-
block.on.buffer.full=true:过时了,可以转而设置max.block.ms -
acks=all:设置 acks 为 all 很容易理解,即必须要等到所有 follower 都响应了发送消息才能认为提交成功,这是 producer 端最强程度的持久化保证。 -
retries=Integer.MAX_VALUE:设置成 MAX_VALUE 纵然有些极端,但其实想表达的是 producer 要开启无限重试 。 用户不必担心 producer 会重试那些肯定无法恢复的错误,当前 producer 只会重试那些可恢复的异常情况,所以放心地设置一个比较大的值通常能很好地保证消息不丢失。 -
max.in.flight.requests.per.connection=1:设置该参数为 1 主要是为工防止 topic 同分区下的消息乱序问题。这个参数的实际效果其实限制了 producer 在单个 broker 连接上能够发送的未响应请求的数量。因此,如果设置成 1,则 producer 在某个 broker 发送响应之前将无法再给该 broker 发送 producer请求 。 -
使用带有回调机制的send
producer.send(new ProducerRecord<>(topic, messageKey, messageStr), new DemoCallBack(startTime, messageKey, messageStr)); //不要使用kafkaProducer中单参数的send方法,该方法仅仅将消息发送出去而不会理会消息发送的结果,如果消息发送失败,该方法不会得到任何通知,可能造成消息的丢失。 //在Callback逻辑中显式立即关闭producer,在 Callback 的失败处理逻辑中显式调用 KafkaProducer.close(0)。这样做的目的是为了处理消息的乱序问题。若不使用 close(0),默认情况下 producer 会被允许将未完成的消息发送出去,这样就有可能造成消息乱序。
-
-
broker端配置
unclean.leader.election.enable=false:关闭unclean.leader选举,不允许非ISR中的副本被选举为leader,从而避免 broker 端因日志水位HW截断而造成的消息丢失 。replication.factor >=3:设置成 3 主要是参考了 Hadoop 及业界通用的三备份原则,其实这里想强调的是一定要使用多个副本来保存分区的消息。min.insync.replicas >1:用于控制某条消息至少被写入到 ISR 中的多少个副本才算成功,设置成大于 1 是为了提升producer 端发送语义的持久性。记住只有在 producer 端 acks 被设置成 all或-1 时,这个参数才有意义。replication.factor > min.insync.replicas:若两者相等,那么只要有一个副本挂掉,分区就无法正常工作,虽然有很高的持久性但可用性被极大地降低了 。 推荐配置成 replication.factor= min.insync.replicas + 1。
-
5.1.7、消息压缩
- 众所周知,数据压缩显著地降低了磁盘占用或带宽占用,从而有效地提升了 IO密集型应用的性能。不过引入压缩同时会消耗额外的 CPU 时钟周期,因此压缩是 IO性能和 CPU 资源的平衡( trade-off) 。
- kafka自0.7x版本便开始支持压缩特性,producer能将一批消息压缩成一条消息发送,而broker端将这条消息写入本地日志文件,consumer拉取这条压缩消息时,它会自动的对消息进行解压缩,还原成初始的消息集合返回给用户。
- 一句话概述:producer端压缩,broker端保持,consumer端解压缩。
- 所谓的broker端保持是通常情况下不会进行解压缩操作,这里的通常情况下需要满足一定的条件,如果一些前置条件(比如需要进行消息的格式转换)不满足时,那么broker端就要对消息进行解压缩然后重新压缩。
5.1.8、压缩算法
- GZIP、Snappy、LZ4
- Zstandard:Zstandard 是少有的能够同时在效率和性能方面超过当前业界“翘楚” Zlib 的压缩算法之一
5.1.9、性能调优
- KatkaProducer. send 方法逻辑的主要耗时都在消息压缩操作上,因此妥善地调优压缩算法至关重要。不过在此之前,首先比较一下目前 Kafka 支持的压缩算法的性能。
- 对 Kafka 而言 ,性能测试的结果出奇地一致,即 LZ4 >> Snappy > GZIP 。 本来 Snappy 和 LZ4 无论在压缩速率还是压缩比率上都差不太多,但由于 Kafka 源代码中对 Snappy 的某个关键参数进行了硬编码,使得 Snappy 与 Kafka 的结合表现并不优秀,至少比 LZ4 要差很多。
- 可以看出kafka对LZ4压缩算法的支持是最好的,启用LZ4进行消息压缩的producer的吞吐量是最高的。
- 如何调优producer的压缩性能:先判断是否启用压缩的依据是 IO 资源消耗与CPU 资源消耗的对 比 。如果生产环境中的 IO资源非常紧张,比如 producer 程序消耗了大量的网络带宽或 broker 端的磁盘占用率非常高,而 producer 端的 CPU 资源非常富裕,那么就可以考虑为 producer 开启消息压缩。反之则不需要设置消息压缩以节省宝贵的 CPU 时钟周期 。

- 其次,压缩的性能与producer端batch的大小息息相关,通常认为batch越大,压缩时间越长。不过时间的增长不是线性的,而是越来越缓慢。如果发现压缩很慢,说明系统的瓶颈在用户主线程而不是I/O发送线程,因此可以考虑增加多个用户线程同时发送消息,这样通常能显著地提升producer吞吐量。
6.2、Consumer
6.2.1、构造consumer实例
-
构造一 个 java.util.Properties 对象,至少指定 bootstrap. servers 、 key.deserializer 、value.deserializer 和 group.id 的值,即上面代码中显式指明必须指定带注释的 4 个参数
-
使用上一步创建的 Properties 实例构造 KafkaConsumer 对象
KafkaConsumer<Integer, String> consumer = new KafkaConsumer<>(props); -
调用 KafkaConsumer.subscribe 方法订阅 consumer group 感兴趣的 topic 列表
consumer.subscribe(Collections.singletonList(this.topic)); -
循环调用 KafkaConsumer. poll 方法获取封装在 ConsumerRecord 的 topic 消息
ConsumerRecords<Integer, String> records = consumer.poll(Duration.ofSeconds(1)); consumer.poll(1000) //要么获取了足够多的可用数据 //要么等待时间超过了指定的超时设置 -
处理获取到的 ConsumerRecord 对象
-
关闭 KafkaConsumer
- 到底哪一步算作 consumer 端的消费?是调用 poll 方法这一步?还是处理 ConsumerRecord 对象这一步?抑或是两者加起来?
- 从kafka consumer的角度讲,poll方法返回即认为consumer成功消费了信息。如果你发现poll返回消息的速度太慢,那么可以调节相应的参数来提升poll方法的效率,若消息的业务及处理逻辑过慢,则应该考虑简化处理逻辑或者把处理逻辑放入单独的线程来执行
6.2.2、Offset
-
consumer 端需要为每个它要读取的分区保存消费进度,即分区中当前最新消费消息的位置 。
该位置就被称为位移( offset ),consumer需要定期向kafka提供自己的位移信息。 -
消息交付语义:
- 最多一次(at most once):消息可能丢失,但不会被重复处理
- 最少一次(at least once):消息不会丢失,但可能会被处理多次
- 精准一次(exactly once):消息一定会处理且只会处理一次
-
消费之前提交位移,即可实现at most once
-
消费之后提交位移,即可实现at least once
-
幂等性,事务实现exactly once
-
consumer 会在 Kafka 集群的所有 broker 中选择 一 个 broker 作为 consumer group 的
coordinator ,用于实现组成员管理、消费分配方案制定以及提交位移等 。 -
consumer 提交位移的主要机制是通过向所属的 coordinator 发送位移提交请求来实现的 。
每个位移提交请求都会往_consumer_offsets 对应分区上追加写入一条消息 。 消息的 key 是
group.id 、 topic 和分区的元组,而 value 就是位移值 。
6.3、rebalance
6.3.1、rebalance概述
- consumer group 的 rebalance 本质上是一组协议,它规定了一个 consumer group 是如何达成
一致来分配订阅 topic 的所有分区的 。 这个分区分配的过程就是rebalance - 对于每个组而言, Kafka 的某个broker 会被选举为组协调者( group coordinator)
6.3.2、rebalance触发条件
- 组成员发生变更,比如新 consumer 加入组,或己有 consumer 主动离开组,再或是己有 consumer 崩溃时则触发 rebalance 。
- 组订阅 topic 数发生变更,比如使用基于正则表达式的订阅,当匹配正则表达式的新topic 被创建时则会触发 rebalance 。
- 组订阅 topic 的分区数发生变更,比如使用命令行脚本增加了订阅 topic 的分区数 。
6.3.3、rebalance分区分配
- range
- round-robin
- sticky
6.3.4、rebalance generation
- 某个 consumer group 可以执行任意次 rebalance。为了更好地隔离每次 rebalance 上的数据,新版本 consumer 设计了 rebalance generation 用于标识某次 rebalance
- 例如,Generation l 时 group 有 3 个成员,随后成员 2 退出组,coordinator 触发 rebalance, consumer group 进入到 Generation 2 时代,之后成员 4 加入,再次触发 rebalance, group 进入到 Generation 3 时代
- 保护消费者组,防止无效 offset 提交 。比如上一届的 consumer 成员由于某些原因延迟提交了 offset,但 rebalance 之后该 group 产生了新一届的 group 成员,而这次延迟的 offset 提交携带的是旧的 generation 信息,因此这次提交会被 consumer group 拒绝 。
6.3.5、rebalance流程
-
consumer group 在执行 rebalance 之前必须首先确定 coordinator 所在的 broker,并创建与该broker 相互通信的 Socket 连接 。 确定 coordinator 的算法与确定 offset 被提交到consumer offsets 目标分区的算法是相同的 。
- 计算 Math.abs(groupID.hashCode) % offsets.topic.num.partitions 参数值(默认是 50) ,假设是 10
- 寻找__consumer_offsets 分区 10 的 leader 副本所在的 broker,该 broker 即为这个 group的 coordinator
-
进行rebalance
-
加入组
-
同步更新

-
6.3.6、rebalance监听器
- 新版本 consumer 默认把位移提交到 consumer offsets 中,Kafka 也支持用户把位移提交到外部存储中,比如数据库中。若要实现这个功能,用户就必须使用 rebalance 监听器 ,使用 rebalance 监听器的前提是用户使用 consumer group 。 如果使用的是独立 consumer 或是直接手动分配分区,那么 rebalance 监听器是无效的
- rebalance 监昕器最常见的用法就是手动提交位移到第三方存储以及在 rebalance 前后执行一些必要的审计操作
6.3.7、多线程消费实例
-
每个线程维护一个kafkaConsumer:每个线程都会创建专属于该线程的kafkaConsumer实例,消费固定数目的分区

-
单个kafkaConsumer实例+多worker线程:我们将消息的获取与消息的处理解耦,把后者放入单独的工作者线程中,即所谓的 worker 线程中,同时在全局维护一个或若干个 consumer 实例执行消息获取任务。
-
用全局的 KafkaConsumer 实例执行消息获取,然后把获取到的消息集合交给钱程池中的 worker 线程执行工作。之后 worker 线程完成处理后上报位移状态,由全局 consumer 提交位移
-
两种方式对比
6.3.8、独立Consumer
-
目前为止我们讨论的consumer都是以consumer group的形式出现的,group自动帮用户执行分区分配和rebalance,对于需要有多个consumer共同读取某个topic的需求来说,使用group是非常方便的。但有时候用户依然有精确控制消费的需求,比如严格控制某个consumer固定消费那些分区:
- 如果进程自己维护分区的状态,那么它就可以固定消费某些分区而不用担心消费状态的丢失
- 如果进程本身已经是高可用且能够自动重启恢复错误(比如使用YARN和Mesos等容器调度框架),那么它就不需要kafka来帮他完成错误检查和状态恢复
-
以上两种情况,consumer group都是无用武之地,取而代之的是独立consumer(standalone consumer)的角色,standalone consumer间彼此独立工作互不干扰,任何一个consumer崩溃都不影响其他standalone consumer的工作。
6.4、broker
- broker承载了绝大多数的kafka服务,对于用户而言,broker的主要功能就是持久化消息以及将消息队列中的消息从发送端传输到消费端。
6.4.1、消息设计
V0版本
- CRC校验码:4字节,保证消息传输过程中不会被篡改
- magic字段:1字节,版本号,V0版本magic=0,V1版本magic=1,V2版本magic=2
- attribute字段:1字节,V0版本,只使用低3位表示压缩格式,其他5位扩展使用
- key长度字段:4字节,表示key字段长度,未指定key,默认-1
- key字段:消息key,长度由“key字段长度”值决定,如果-1,值为0
- value长度字段:4字节,未指定值为-1
- value字段:跟“key字段”相同
- 除了key和value字段之外的所有字段统称为“消息头部信息(message header)”,占14字节,也就是说V0版本的消息长度最小为14字节,否则就是非法消息
V1版本
-
V0版本的缺陷
- 没有消息的时间信息,kafka消息过期日志只能依靠日志段文件的“最近修改时间”,但这个时间极易受到外部操作的干扰。举个例子,若对日志文件执行touch命令,则该日志文件的修改时间就更新了,一旦时间被修改或者被破坏,kafka无法正确判断哪些为过期消息
- 许多流式框架都需要使用消息保存时间信息对消息执行时间窗口等聚合操作
-
V1版本引入了8字节的时间戳字段,另外attribute字段的第4位用于保存时间戳类型,其他字段与V0版本含义相同。
-
当前支持两种时间戳类型:CREATE_TIME和LOG_APPEND_TIME,前者有producer指定,后者表示消息发到broker端由broker指定
-
V0版本和V1版本的缺陷
- 空间利用率不高
- 只保存最新位移信息
- 冗余的CRC校验:若用户指定时间戳类型是 LOG APPEND TIME,broker 将使用当前时间戳覆盖掉消息己有时间戳,那么当 broker 端对消息进行时间戳更新后, CRC 就 需要重新计算从而发生变化;再如, broker 端进行消息格式转换( broker 端和 clients 端要求版本不一致时会发生消息格式转换,不过这对用户而言是完全透明的)也会带来 CRC 值的变化 。 鉴于这些情况,对每条消息都执行 CRC 校验
实际上没有必要,不仅浪费空间,还占用了宝贵的 CPU 时间片 - 未保存消息长度
V2版本
- length:消息总长度
- arribute:弃用,但是在消息格式中仍占1字节,扩展使用
- timestamp delta:时间增量,这里保存与RecordBatch的其实时间戳的差值的话可以进一步的节省占用的字节数
- offset delta:位移增量,保存与RecordBatch起始位移的差值,可以节省占用的字节数
- headers:
- CRC校验:在v2版本中将crc的字段从Record中转移到了RecordBatch中
- record batch
- first offset:表示当前RecordBatch的起始位移
- first timestamp:RecordBatch中第一条Record的时间戳