kafka讲解
kafka定义:
kafka是一个高吞吐,分布式的发布/订阅消息系统。基于日志存储; producer(生产者):负责broker推送消息;
consumer(消费者):负责从broker拉取消息消费; Consumer Group:消费者组则是存在多个consumer,
消费者消费broker中当前Topic的不同分区中的消息,消费者组之间互不影响; brocker:kafka集群中的单台服务器;
topic(主题):主题,生产者和消费者通过topic来对接 partition 分区:在物理上对应 Broker 磁盘上的一个独立目录
Replica (副本):容灾备份。分为 Leader (读写) 和 Follower (同步),防止节点故障导致数据丢失;
应用场景:
解耦:消息队列可以作为一个接口层,解耦重要业务流程
限流削峰:使用消息队列作为中间件,可以将流量的高峰保存在消息队列中,从而防止系统的高请求,减轻服务器的请求压力。
异步处理:消息队列提供了异步处理机制,允许用户把一个消息放入到队列中,但并不立即处理它,然后再需要的时候进行处理。
Kafka 的写入流程:
Producer 首先根据 Key Hash 取模(或轮询)确定目标 Partition,将消息在本地缓冲区累积成 Batch
批次,随后推送给该 Partition 的 Leader Broker。 Leader 接收后,将消息顺序追加写入本地磁盘
Segment,同时 ISR 列表中的 Follower 会主动拉取并同步数据。 在 acks=all 模式下,只有当 Leader 和所有
ISR Follower 都确认写入成功后,Leader 才会向 Producer 返回 ACK。 若 Producer 未收到
ACK(如网络超时或 Leader 宕机),会自动触发重试机制(重试前会刷新元数据定位新 Leader),并依靠幂等性保证消息不重不漏。
整个流程通过批量发送和顺序写盘保证了高吞吐,通过多副本同步和ACK 机制确保了数据零丢失。
Kafka 的消费流程:
Consumer 启动后首先加入指定的 Consumer Group,触发 Rebalance(重平衡) 机制,由 Group
Coordinator (卡儿德奶德)(集群协调节点)将 Topic 下的所有 Partition 均匀分配 给组内消费者,确保每个
Partition 仅被组内一个 Consumer 独占。 随后,Consumer 主动拉取(Pull) 数据,直接向 Partition
的 Leader Broker 发送请求并指定 Offset(消费位移), Broker 返回消息批次后,Consumer
执行本地业务逻辑。 处理完成后,Consumer 将最新的 Offset 提交(Commit) 到 Kafka 内部的
__consumer_offsets 主题中(推荐手动提交以确保“至少一次”语义)。 若消费者宕机重启,会从已提交的 Offset 处 断点续传;若消费过程中发生 重平衡,未提交 Offset 的消息会被重新消费,从而在保证高并发扩展性的同时, 实现数据的可靠处理与不丢失。
Kafka消费模式:
一对一的消费,也就是点对点的消费,一个发送一个接收,发送消息到消息队列,消费者根据消息队列的订阅进行拉取消息消费。
一对多,发布者将消息发送到topic中,系统将这些消息传递给多个订阅者。
Kafka中的ISR(InSync-Repli)、OSR(OutSyncRepli)、AR(AllRepli)又代表什么?
Kafka 通过 AR 管理所有副本,其中跟得上 Leader 进度的称为 ISR,落后的称为 OSR。 只有 ISR
中的副本才有资格选举为新 Leader,保证高可用。AR(分配的副本):分区下所有的副本集合(包含正常的和失败的)
ISR(同步副本):kafka中follower的消息基本能与Leader保持一致,如果Leader出现故障,新的leader将从isr中选出,若follower落后太多则会被踢出到ORS中
ORS(不同步队列副本):被移除的ISR集合,追赶数据后自动回归ISR;
失效副本是指什么:
失效的副本为速率比leader相差大于10s的follower,将失效的副本现先剔除出ISR中,等速率接近leader10s在加入到ISR中。
Kafka中有哪些需要选举的地方?
在ISR中需要选举,选择策略为先到先得。
Kafka创建topic的时候如何将分区放入不通的Broker中:
首先副本数不超过Broker数量,第一个分区是随机从Broker中选择一个,然后其他分区相对于0号分区一次向后移。
Kafka中HW、LEO等分别代表什么:
LEO (Log End Offset,日志末端偏移):每条副本日志中下一条待写入消息的位置(即当前最后一条消息的 offset +
1),每个副本都有自己的 LEO,Leader 的 LEO 通常最大。
HW(High Watermark,高水位):ISR 集合中所有副本 LEO 的最小值。消费者只能读取到HW之前的消息
Kafka消息采用的是pull模式,还是push模式?
生产者 -> Broker: Push (推) 模式。生产者主动将消息发送给 Broker。 Broker -> 消费者: Pull (拉)
模式。消费者主动向 Broker 请求数据。
Kafka中如何体现消息顺序性:
每个分区内,每条消息都有offset,所以只能在同一分区内有序,但不通的分区无法做到消息顺序性。
Kafka分区的目的:
负载均衡,对于消费者来说可以提高 高并发,提高读取效率。
消费者提交消费位移时提交的是当前消费的最新消息的offer还是offset +1:
生产者发送数据offer是从0开始的,消费者消费的数据offset是从offset +1 开始的
Kafka中的分区器、序列化、拦截器他们之间的处理顺序是什么:
拦截器》序列化器》分区器
Zk在kafka中的作用:
集群元数据管理:kafka集群中所有的Broker信息,topic和分区的状态信息都会储存在zk节点上。
负责进行leader选举:kafka中每一个分区都会有一个leader,zk负责进行leader的选举,当一个leader
dead掉后,集群可以快速选择新的leader继续服务。
kafka选举:
zk传统模式:所有brock启动的时候,都会zk的/controller节点注册watch.第一个启动的brocker成功创建临时节点,成为controller.其他创建失败的,则进入监听状态,当controller宕机的时候,重新竞争选出新的controoler;
新版本选举:
KRafr(新版本):基于KRafr协议,controoller节点内部通过Raft选举,产生新的leader,不在需要zk
分区leader选举:
由controller 负责发起,当分区的leader宕机,controller从该分区的ISR中选举一个出来
消费组选主:
如何消费者内还没有leader,那么第一个加入消费组的消费者即为消费者的leader。
消费端如何保障数据不丢失:
采用消费者的offse偏移量,监听消费的topic,从kafka中获取上一次消费到的那个偏移,开始消费,当消费完成后,需要向kafka报告消费完成更新偏移量信息。
生产者如何保障数据不丢失:
让生产者在生产数据到broker端,可以通过设置ack校验,让broker给予ack响应,从而判断数据是否已经写入broker端,设置ack主要有三种:
(1)acks=0 生产之只管发送数据,不关心broker是否已经接收数据。 (2)acks=1 (默认) 生产者将数据发送到broker,
只等待 Leader (主副本) 写入成功。立刻返回 ACK。 (3)acks=all 生产者将数据发生到broker,等待
Leader + 所有 ISR (同步副本) 写入成功。,返回ack,才认为写入成功。
数据发送一条,响应一次ack,如果broker迟迟不给响应如何解决?
1.重试策略 配置 retries 参数,修改重试次数,如果重试3次依旧没有响应 则报错。
2.开启幂等性生产者 (enable.idempotence=true)。这样即使因为超时而重试,Kafka 也能通过序列号在 Broker 端去重,保证数据不重复。
原理图:
kafka 数据存储机制

kafka架构说明:

kafka数据分发策略:

kafka的消费者负载均衡策略:

分片和副本机制:

如何保证生产端数据不丢失:

如何保证消费端数据不丢失:

更多推荐
所有评论(0)