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的消费者负载均衡策略:
在这里插入图片描述
分片和副本机制:
在这里插入图片描述
如何保证生产端数据不丢失:
在这里插入图片描述
如何保证消费端数据不丢失:
在这里插入图片描述

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐