ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

kafka的文章

kafka的文章 浅谈kafka从Kafka中学习高性能系统如何设计 | 京东云技术团队1.面试的问题要点 至多一次、恰好一次数据一致性超时重试、幂等消息顺序消息挤压延时消息1.1 kafaka 生产消息的过程。在消息发送的过程中涉及到了两个线程一个是main 线程一个是sender 线程。在main 线程中创建了一个双端队列 RecordAccumulatormain 线程将消息发送到 双端队列sender 线程不断从双端队列读取 发送到 broker1.2 消息队列的可靠性。1.3 副本同步机制leo: 定义LEO 即日志末端偏移量它表示每个副本日志中最后一条消息的下一个偏移量。hw:高水位HWHigh Watermark的确是 ISRIn - Sync Replicas同步副本集合中所有副本的最小日志末端偏移量LEOLog End OffsetKafka 的副本同步机制Kafka 的副本同步机制是保障数据可靠性和高可用性的核心特性下面从整体架构、同步流程、ISR 机制、相关参数等方面进行详细介绍。整体架构Kafka 中每个分区都有一个 Leader 副本和多个 Follower 副本。生产者和消费者只与 Leader 副本进行交互Follower 副本负责从 Leader 副本同步数据。这样的设计使得 Kafka 可以在多个 Broker 上存储数据副本提高数据的容错能力。同步流程消息生产生产者将消息发送到 Kafka 集群时会指定要发送到的主题和分区。Kafka 根据分区的 Leader 副本位置将消息发送到对应的 Leader 副本所在的 Broker。Leader 副本接收消息Leader 副本接收到生产者发送的消息后将消息写入本地日志并更新自身的日志末端偏移量LEO。Follower 副本同步消息Follower 副本通过向 Leader 副本发送 Fetch 请求来同步消息。Fetch 请求中包含 Follower 副本当前的 LEOLeader 副本根据该信息将新的消息发送给 Follower 副本。Follower 副本接收到消息后将消息写入本地日志并更新自身的 LEO。高水位HW更新高水位HW是分区中所有副本都已经成功复制的消息的最大偏移量。Kafka 会根据 ISRIn - Sync Replicas同步副本集合中所有副本的 LEO 来更新 HW。具体来说HW 是 ISR 中最小的 LEO。只有偏移量小于 HW 的消息才被认为是已经在所有同步副本中安全保存的可以被消费者消费。ISR 机制ISR 定义ISR 是与 Leader 副本保持同步的一组副本集合。只有在 ISR 中的副本才被认为是可靠的同步副本能够参与 HW 的计算。ISR 动态维护Kafka 会定期检查 Follower 副本与 Leader 副本的同步情况通过比较 LEO 的差距来判断 Follower 副本是否落后。如果 Follower 副本的 LEO 与 Leader 副本的 LEO 差距超过一定阈值由 replica.lag.time.max.ms 参数控制则该 Follower 副本会被从 ISR 中移除。当落后的 Follower 副本追上 Leader 副本后它可以重新加入 ISR。ISR 的作用提高数据可靠性只有 ISR 中的副本参与 HW 的计算确保消费者只能读取到已经在多个副本中安全保存的消息。故障转移当 Leader 副本出现故障时Kafka 会从 ISR 中选举出新的 Leader 副本保证数据的一致性和服务的连续性。相关参数acks 参数该参数用于控制生产者发送消息时的确认机制影响副本同步策略。acks 0生产者发送消息后不等待任何确认相当于异步复制性能最高但数据可靠性最低。acks 1生产者发送消息后等待 Leader 副本确认只要 Leader 副本写入成功就返回响应性能和可靠性适中。acks all 或 acks -1生产者发送消息后等待所有 ISR 中的副本确认相当于同步复制数据可靠性最高但性能最低。min.insync.replicas 参数用于指定 ISR 中最少需要有多少个副本同步消息才能认为消息写入成功。结合 acks all 使用时可以进一步增强数据的可靠性。如果 ISR 中的副本数量小于 min.insync.replicas生产者发送消息时会收到写入失败的响应。replica.lag.time.max.ms 参数该参数定义了 Follower 副本与 Leader 副本之间允许的最大延迟时间。如果 Follower 副本在该时间内没有向 Leader 副本发送 Fetch 请求或者没有追上 Leader 副本的 LEO则会被从 ISR 中移除。异常情况处理Leader 副本故障当 Leader 副本所在的 Broker 出现故障时Kafka 会从 ISR 中选举出新的 Leader 副本。新的 Leader 副本会将 HW 作为新的起始偏移量继续处理生产者和消费者的请求。Follower 副本故障如果某个 Follower 副本出现故障它会被从 ISR 中移除。当该副本恢复正常后会重新向 Leader 副本发送 Fetch 请求追赶 Leader 副本的进度当追上后可以重新加入 ISR。综上所述Kafka 的副本同步机制通过 Leader - Follower 架构、ISR 机制和相关参数的配置在保证数据可靠性和高可用性的同时兼顾了性能和容错能力。1.4 kafka 高性能、高吞吐原因磁盘顺序读写顺序读 会使用预读保证了消息的堆积 相比于内存。使用了零拷贝的技术分区分段 索引每个 分区 在磁盘上 按照segment 文件存储的。针对segment 建立.index的索引文件批量压缩 多条消息批量压缩传输降低带宽批量读写1.5 消息丢失的场景 解决方案ackall,配置 min.insync.replicas11.6 消息可靠性的解决方案消息发送ack -1/all 、unclean.leader.election.enable: false, 禁止选举 isr 以外的follower为leadertries 1 重试次数min.insync.replicas1 同步副本数没满足该之前不提供 读写服务。综上所述在 acks all 且 min.insync.replicas 3副本总数为 5 个的情况下至少 3 个处于 ISR 中的副本写入数据完成Kafka 才会判定消息写入操作完成。消费者手动提交 offsetbroker 减少刷盘间隔事务消息1.7 kafka reblance消费者分区策略range 范围分区 默认roundrobin 轮询sticky 策略 体现在 reblance 策略下。触发reblance 的时间消费者组成员个数变化的时候。有新的消费者加入、离开消费者组订阅的topic 发生变化订阅topic 的分区发生变化coordinator 协调过程消费者 找到消费者组中的 协调器确定分区策略1.8 kafka 事务面试官Kafka 事务是如何工作的Kafka的消息会丢失和重复吗——如何实现Kafka精确传递一次语义system design interview 书第八章 设计短链系统SystemDesign系统设计访谈-业内人事指南 《System Design Interview-An insider‘s guide》全书中文翻译收费的几章去这里看面试场景题ddia《DDIA 逐章精读》小册kafka 笔记尚硅谷-笔记面试笔记 kafkaKafka 底层原理与源码分析实战kafka 查询log文件过程在 partition 中通过 offset 查找 message过程根据 offset 的值查找 segment 段中的 index 索引文件。由于索引文件命名是以上一个文件的最后一个offset 进行命名的所以使用二分查找算法能够根据offset 快速定位到指定的索引文件找到索引文件后根据 offset 进行定位找到索引文件中的匹配范围的偏移量position。kafka 采用稀疏索引的方式来提高查找性能得到 position 以后再到对应的 log 文件中从 position处开始查找 offset 对应的消息将每条消息的 offset 与目标 offset 进行比较直到找到消息比如说我们要查找 offset2490 这条消息那么先找到00000000000000000000.index, 然后找到[2487,49111]这个索引再到 log 文件中根据 49111 这个 position 开始查找比较每条消息的 offset 是否大于等于 2490。最后查找到对应的消息以后返回面试题Kafka都有哪些组件controller 选择Controller依赖ZooKeeper实现Controller选举主要是借助于/controller临时节点和ZooKeeper的监听器机制。Controller触发场景有3种集群启动时/controller节点被删除时/controller节点数据变更时。最先 在zk 上创建临时节点/controller 成功的broker 就是controller。最先在Zookeeper上创建临时节点/controller成功的Broker就是Controller。Controller重度依赖Zookeeper依赖zookeepr保存元数据依赖zookeeper进行服务发现。Controller大量使用Watch功能实现对集群的协调管理。如果此时作为Controller的Broker节点宕掉了。那么zookeeper的临时节点/controller就会因为会话超时而自动删除。而监控这个节点的Broker就会收到通知而向ZooKeeper发出创建/controller节点的申请一旦创建成功那么创建成功的Broker节点就成为了新的Controller。脑裂的解决kafka 通过 一个 epoch 纪元的概念。每一个Broker当选Controller时会告诉当前Broker是第几任Controller一旦重新选举时这个任期会自动增1那么不同任期的Controller的epoch值是不同的那么旧的controller一旦发现集群中有新任controller的时候那么它就会完成退出操作清空缓存中断和broker的连接并重新加载最新的缓存让自己重新变成一个普通的Broker。Producer生产消息是如何实现的 引导[主线程] 调用 send()│▼[拦截器] → [序列化器] → [分区器] → [RecordAccumulator内存缓冲区]│▼[Sender 线程]后台独立线程定时唤醒│▼[NetworkClient]NIO 多路复用 → [Broker Leader]│▼响应回调 / 重试 / 元数据更新一、三层回答法由浅入深层层递进第1层一句话概括总体流程10秒定调先给出高度凝练的一句话展示你对全局的把控同时给后续展开留出空间。“Kafka 生产者的消息发送本质是一个异步双缓冲 批量 NIO 网络模型的过程。消息经过拦截、序列化、分区后先暂存到内存累积器再被后台发送线程按 Broker 聚合批量发出。”第2层按步骤讲清楚核心链路1-2分钟展示细节这里可以顺着链路有详有略地讲重点突出缓冲区和发送线程的协作因为这最能体现你对内部机制的理解。建议这样展开拦截器非必须但不提不扣分提了显得你全面。一笔带过即可。序列化将 Key/Value 转为字节数组略讲。分区选择三种情况指定分区 / 有 Key 哈希 / 无 Key 粘性重点提一下“粘性分区”是为了把消息填满一个 batch 再换分区提升批量效率——这里可以埋第一个钩子见后文。RecordAccumulator这是重头戏。你要讲清楚它是一个“按分区组织的双端队列每个队列里装着 ProducerBatch”新消息优先追加到队列尾部的未满 batch满了就新建 batch。这样设计是为了拼批。Sender 线程后台的 NIO 线程会把多个分区的 batch 按目标 Broker 节点重新聚合打包成一个请求异步发出。顺口提一句 max.in.flight.requests.per.connection 控制并发请求数为后面讲顺序性埋钩子。Broker 响应与重试根据 acks 配置返回响应可重试错误会退避重试幂等性下能消除重复。第3层主动抛出你准备好的“亮点钩子”这是引导面试官的关键讲完基本流程后不要等面试官来问你可以主动说一句过渡语把话题引向你擅长的领域“以上是基础的发送流程但其实这里面有几个点非常值得深入比如如何保证发送既不丢又不重、如何平衡延迟和吞吐、以及幂等性和事务的实现原理我可以挑一个您感兴趣的展开聊聊吗”如果面试官让你自己选你当然选最熟的那个很多时候面试官会直接顺着你的钩子提问你就完全掌握了主动权。二、四个可以主动引导的“钩子”及话术钩子1可靠性保障ack 与幂等性—— 最容易出彩在你讲完流程后可以这样铺垫“其实Producer 设计的精妙之处在于它通过 acks 重试 幂等性 三个机制的组合几乎可以做到精确一次语义。”如果面试官追问你就展开acksall 保证数据多副本确认不丢。重试机制可能造成“发了一次网断重发导致重复”幂等性通过 PID 分区级别的序列号让 Broker 识别并丢弃重复消息完美解决。再进阶到事务跨分区原子性写入提及 initTransactions 等 API 和应用场景。钩子2批量与吞吐优化 —— 展示性能调优能力讲 RecordAccumulator 时自然带出“这里其实牵涉到一个经典的吞吐与延迟取舍batch.size 和 linger.ms 就是用来控制拼批力度的。我比较推崇的做法是把 linger.ms 从 0 调到 5-100ms这样业务几乎感知不到延迟但吞吐量能成倍提升因为网络请求大幅减少。”面试官如果感兴趣就会追着问你具体的调优经验你就可以展开讲压缩算法选择、缓冲区大小设置等。钩子3顺序性保证 —— 展示你对并发和幂等性的深度理解在提 Sender 线程时埋钩子“另外要保证单分区有序传统做法是把 max.in.flight.requests.per.connection 设为 1但这会牺牲吞吐。好在开启幂等性后Broker 会按序列号排序写即使并发发送也不会乱序这个设计非常巧妙。”如果你主动提到这点面试官大概率会问你“为什么幂等性可以保证顺序”你就可以详细解释 PID 和 sequence number 在服务端的处理逻辑。钩子4分区选择与数据倾斜 —— 展示方案设计能力在讲 Partitioner 时“分区路由里有一个容易被忽视的坑默认用 Key 哈希时如果 Key 分布不均会导致严重的分区数据倾斜。我在实际项目中遇到过后来是通过自定义分区器或对 Key 做二次加工解决的。另外粘性分区也很值得聊聊它在没有 Key 时比轮询聪明得多……”这能展示你有实际踩坑和解决问题的经验不是纸上谈兵。三、完整回答示例可背诵的流畅话术如果需要一段直接能用的回答你可以这么说开局总览一个 send() 调用背后是一条“异步批量发送”的链路。消息在客户端先写入内存缓冲区由后台线程异步批量发给 Broker。链路细节具体来说消息会先经过可选的拦截器然后被序列化成字节数组接着通过分区器决定发往哪个分区——如果有 Key 就哈希取模没有 Key 就用粘性分区让一批消息尽量粘在一个分区上把 batch 填满。选好分区后消息被追加到 RecordAccumulator 里。这是一个按分区组织的双端队列每个分区队列里装着多个 ProducerBatch。新消息优先追加到队尾未满的批次中满了就建新批次。当某个批次因大小达到 batch.size 或因 linger.ms 时间到期就会被关闭并移交给后台的 Sender 线程。Sender 线程基于 NIO会把这些批次按目标 Broker 节点聚合成一个 ProduceRequest通过一个连接异步发出。Broker 会根据 acks 配置决定何时返回响应acksall 时 Leader 会等待所有 ISR 副本写入完成才确认。如果发生可重试错误生产端会在 retries 内退避重试。抛出钩子引导在这个基本流程之上有几个点设计得非常经典比如幂等性 事务是如何在重试基础上实现精确一次语义的以及粘性分区和批量参数到底如何影响端到端延迟与吞吐。如果您感兴趣我可以选一个深入展开。Follower拉取Leader消息是如何实现的针对“Follower拉取Leader消息是如何实现的”这道题面试官想听到的不仅是“谁拉谁”更是拉取过程中的水位线推进机制、一致性保障和性能设计。和之前一样我们用一条主线 层层递进 埋钩子的方式回答。面试回答可以直接讲总体概述Kafka 的副本同步采用的是Follower 主动向 Leader 拉取Pull 的模式而不是 Leader 推送。这个拉取过程的核心是围绕着 LEO日志末端偏移量 和 HW高水位 这两个关键指标来协同推进的。第一步发出拉取请求每个 Follower 内部有一个独立的同步线程它会持续地向 Leader 发送 Fetch 请求。这个请求里最关键的一个参数就是 Follower 自己当前的 FetchOffset——也就是它本地下一条要写入消息的位置。通过这个值Leader 就能知道这个 Follower 的复制进度即 Follower 的 LEO。第二步Leader 处理并返回Leader 收到请求后会根据 FetchOffset 去本地的日志段中读取数据。这里有一个很大的性能优化零拷贝。Leader 会利用操作系统的 sendfile 系统调用直接将磁盘上的数据从内核态缓存传输到网卡绕过应用层内存最大限度地减少 CPU 拷贝开销。在返回给 Follower 的响应中除了拉取到的消息还包含了一个非常重要的字段Leader 当前的 HW高水位。这个 HW 告诉 Follower“在这个偏移量之前的消息都已经被所有 ISR 副本确认可以被消费者消费了。”第三步Follower 写入并推进 HWFollower 收到这批消息后先写到本地的日志文件中同时更新自己的 LEO。然后它会用自己新的 LEO 作为下一次 Fetch 请求的 FetchOffset。同时Follower 会拿 Leader 返回的 HW 和自己的 LEO 做一个比较取两者中的较小值来更新自己的 HW。这样Follower 上的 HW 就能安全地向前移动保证不会超出 Leader 认为已提交的范围。第四步Leader 侧计算并推进 HW这个拉取过程反过来也推动了 Leader 端的 HW 计算。Leader 会实时跟踪所有处于 ISR 列表中的副本的 LEO这个 LEO 就是每个 Follower 在 Fetch 请求中携带的 FetchOffset。Leader 会取所有这些 ISR 副本 LEO 中的最小值作为该分区最新的 HW。一旦这个最小值增大Leader 的 HW 就向前推进意味着有新的消息正式变为“已提交”可以被消费者读取了。小结所以Follower 拉取的过程本质上就是一次请求双向推进Follower 通过拉取消息推进自己的 LEOLeader 通过 Follower 请求中的 FetchOffset 更新该 Follower 的 LEO进而重新计算 HWFollower 再通过下一次响应的 HW 更新自己的 HW形成闭环。面试中如何引导你的钩子策略这个机制里藏着很多可以展示你深度的“钩子”讲完上述流程后你可以主动说一句“上面就是基于 HW 的拉取同步基本流程。不过这个原始设计其实有两个非常值得深入的点一个是HW 可能带来的数据丢失或数据不一致风险以及 Kafka 后来如何通过 Leader Epoch 机制来解决另一个是ISR 的动态管理比如 Follower 因滞后被踢出后HW 如何重算。您想听我展开聊聊哪一个”这样一来主动权就在你手里了。常用的引导钩子有钩子1HW 的缺陷与 Leader Epoch 的救场最显深度你可以讲在发生 Leader 切换时新 Leader 的 HW 可能比旧 Leader 更保守导致部分已提交消息“回退”丢失或者 Follower 重启后无脑截断到 HW可能导致日志不一致。而 Leader Epoch 机制让每个 Leader 周期带上版本号Follower 恢复时只从上一个 Epoch 的确认位置之后截断精准又安全。钩子2ISR 的动态维护展示系统思维你可以讲Follower 如果长时间不拉取或者拉取偏移量落后太多Leader 就会把它从 ISR 中移除。这时min.insync.replicas 配置和 HW 的推进都会受影响。这直接关系到 acksall 时的可用性和可靠性权衡。·min.insync.replicas是 Kafka 中控制数据可靠性与可用性权衡的关键参数它定义了生产者请求确认acks时必须成功写入的最小副本数。钩子3为什么不采用 Push 模式展示架构理解你可以讲Kafka 选择 Pull 模式是为了让 Leader 无状态化Follower 自己控制拉取节奏天然支持磁盘慢、网络慢的异构集群。而且这和消费者端的拉取模型是一致的简化了整体架构。钩子4零拷贝的具体实现展示操作系统功底你可以讲Leader 读取日志时如果使用普通的 read/write 会有四次上下文切换和两次 CPU 拷贝而通过 Java 的 FileChannel.transferTo() 调用的 sendfile只需要两次上下文切换和一次DMA拷贝吞吐量能提升几十倍。leader 宕机后 选哪个副本从 isr 中选取第一个follower作为leader[Leader,follower-1,follower-2]数据同步一致性问题 hw 截断Consumer拉取消息是如何实现的针对“Consumer 拉取消息是如何实现的”这道面试题回答的重点在于把消费者组协调机制、分区分配、拉取循环、位移管理这四块串成一个闭环并展示你对可靠性与一致性权衡的理解。下面按面试场景给出回答与引导。面试回答可以直接讲总体概述Kafka 消费者的核心是基于拉取Pull模式的分布式消费模型由消费者组、GroupCoordinator、分区分配策略三者协同工作。消费者会先找到协调器、加入消费组并分配到分区然后在一个循环中向 Broker 发送 Fetch 请求拉取消息最后提交消费位移实现进度管理。第一步寻找 GroupCoordinator消费者启动时会配置 group.id。它会根据 group.id 的哈希值找到 __consumer_offsets 主题的对应分区该分区的 Leader Broker 就是当前消费者组的 GroupCoordinator。后续所有组管理交互加入组、同步组、心跳、位移提交都通过这个协调器完成。第二步加入消费者组并分配分区消费者找到协调器后会经历两个核心请求JoinGroup 请求消费者向协调器报告自己的订阅信息协调器会选出组内第一个加入的消费者作为 Group Leader。SyncGroup 请求Group Leader 根据配置的分区分配策略Range、RoundRobin、Sticky、CooperativeSticky计算出分区分配方案通过 SyncGroup 请求发给协调器。协调器再将该方案通知给组内所有消费者。每个消费者最终知道自己被分配了哪些分区。第三步拉取循环核心拿到分区后消费者会进入一个 while 循环不断调用 poll()。其内部运作是发送 Fetch 请求消费者向各分区的 Leader Broker 发送 Fetch 请求请求中携带自身已经消费到的 fetch offset对应于下一条要拉取消息的位置。零拷贝与数据返回Broker 收到后通过 sendfile 零拷贝技术直接将磁盘日志读取并传输给消费者响应中包含了消息集。客户端管理客户端维护一个 Fetcher 组件它会异步预读取数据消费者网络层poll() 只是从内部队列中取出已准备好的数据从而实现拉取的“伪同步”高效模式。位移更新poll() 返回消息后消费者可以异步或同步地提交位移offset。默认是自动提交enable.auto.committrue但更常见的生产级做法是手动提交以保证处理完成才推进位移。第四步心跳与再均衡消费者会在后台定期发送 Heartbeat 请求给 GroupCoordinator以维持自己在组内的活状态。一旦出现消费者宕机、网络分区或新消费者加入协调器会触发 Rebalance所有消费者被要求重新加入组Revoke Join暂停消费。分区分配方案被重新计算并分配。消费者拿到新分区后再次进入拉取循环。小结所以消费者拉取消息的整体链路就是找协调器 → 加入组拿分区 → 以 fetch offset 为参数不断拉取 → 按需提交位移而协调器通过心跳和再均衡保证了消费组的动态伸缩和高可用。面试中如何引导你的钩子策略讲完以上流程后可以主动引导一句“以上就是消费者拉取的基本流程。不过这里面有五个非常值得深入的点可以直接反映工程师对 Kafka 消费者的理解和经验手动位移提交的语义控制、再均衡的问题与优化、为什么消费性能与分区数强相关、消费者拉取如何做到顺序性以及 Kafka 为什么坚持 Pull 模型。您感兴趣的话我可以挑任何一个展开聊聊。”当你掌握主动权后可根据面试官兴趣选择以下钩子深入钩子1位移提交与消费语义解释 自动提交 导致 at-most-once 或消息重复的问题。手动同步/异步提交 如何实现 at-least-once。将位移与业务结果绑定存储如放入事务实现 exactly-once 的思路。钩子2再均衡Rebalance的痛点与优化提及 “消费停顿”问题传统 Eager Rebalance 会 Stop-The-World所有消费者放弃分区再重新分配。引入 Cooperative Sticky Rebalance 后分区可以逐步再分配减少了不必要的停顿。建议缩小心跳间隔、增大 session.timeout 等调优经验。钩子3消费者数量与分区数的强关联性强调“消费者的并行度上限 主题分区数”。多余消费者会空闲反之单个消费者处理多个分区。可引出如何通过分区对齐和业务分区策略设计高吞吐消费方案。钩子4消费的顺序性保证Kafka 只保证单分区有序消费者端顺序取决于分区分配和消费逻辑。如何处理多分区下“全局有序”需求通常让所有消息路由到一个分区。钩子5为什么消费者也用 Pull 模型与 Follower 拉取逻辑一脉相承避免 Broker 端推送导致消费者被压垮。消费者可根据自己的处理能力控制拉取速度和批量大小max.poll.records 等实现自然背压。可以与 KafkaConsumer 的暂停/恢复 API 结合展示你对背压控制的掌握。0拷贝什么时0拷贝-小林coding
返回列表