恒美微站 Logo 恒美微站
  • 首页
  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心
  • 联系我们

【消息队列】原理初探之Kafka

  • 首页
  • 资讯中心
  • /
  • 【消息队列】原理初探之Kafka

相关资讯

Hermes Agent 权限边界实战:授权、白名单、命令审批搭起生产安全线 2026/10/9 18:14:13
基于PCA9422与PIC24HJ的便携设备电源管理设计实战 2026/10/9 18:14:13
【ChatGPT】GitHub Copilot 免费注册及在 VS Code 中的安装使用:TaoToken 统一 Key 配置与验证 2026/10/9 18:14:13

最新资讯

数据结构电梯模拟毕业设计:SCAN调度算法与状态机实现
SQL Server 2000 备份还原实战:三类备份策略与四层校验链
Oracle 11.2.0.4 PSU补丁深度解析:修复共享池争用与ORA-00600核心故障
Android Studio推箱子Java游戏开发:地图逻辑与自定义View实现
PyQt5酒店管理系统课设实战:含ER图、SQL建库与PDF报表
高考满分作文与数学45分:偏科生的天赋困境与突围路径

今日推荐

AI编程智能体实战:从写代码到指挥代码的架构与落地
多模态大模型全栈能力拆解:从数据对齐到弹性推理
大模型Agent开发入门:从工具调用循环到落地避坑指南

本周热门

MR25H40CDF + PIC18F65K40:工业记录仪高可靠存储实战
基于STM32的数控恒压恒流电源设计:从硬件到PID调参全解析
LT9211 MIPI重定时器原理与双路扇出实战指南

本月精选

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证
2026 大模型集体涨价:用 Python 做企业 Token 成本测算与选型避坑(附配置)

【消息队列】原理初探之Kafka

发布时间:2026/10/9 18:14:13
【消息队列】原理初探之Kafka Kafka 是由Linkedin公司开发的它是一个分布式的支持多分区、多副本基于Zookeeper的分布式消息流平台它同时也是一款开源的基于发布订阅模式的消息引擎系统。2.1 基本概念消息Kafka中的数据单元被称为消息也被称为记录可以把它看作数据库表中某一行的记录。批次为了提高效率 消息会分批次写入Kafka批次就代指的是一组消息。主题消息的种类称为 主题Topic可以说一个主题代表了一类消息相当于是对消息进行分类。主题就像是数据库中的表。分区主题可以被分为若干个分区partition同一个主题中的分区可以不在一个机器上有可能会部署在多个机器上由此来实现 kafka 的伸缩性单一主题中的分区有序但是无法保证主题中所有的分区有序。生产者向主题发布消息的客户端应用程序称为生产者Producer生产者用于持续不断的向某个主题发送消息。消费者订阅主题消息的客户端程序称为消费者Consumer消费者用于处理生产者产生的消息。消费者群组生产者与消费者的关系就如同餐厅中的厨师和顾客之间的关系一样一个厨师对应多个顾客也就是一个生产者对应多个消费者消费者群组Consumer Group指的就是由一个或多个消费者组成的群体。偏移量偏移量Consumer Offset是一种元数据它是一个不断递增的整数值用来记录消费者发生重平衡时的位置以便用来恢复数据。broker: 一个独立的 Kafka 服务器就被称为 brokerbroker 接收来自生产者的消息为消息设置偏移量并提交消息到磁盘保存。broker 集群broker 是集群 的组成部分broker 集群由一个或多个 broker 组成每个集群都有一个 broker 同时充当了集群控制器的角色自动从集群的活跃成员中选举出来。副本Kafka中消息的备份又叫做 副本Replica副本的数量是可以配置的Kafka 定义了两类副本领导者副本Leader Replica 和 追随者副本Follower Replica前者对外提供服务后者只是被动跟随。重平衡Rebalance。消费者组内某个消费者实例挂掉后其他消费者实例自动重新分配订阅主题分区的过程。Rebalance是 Kafka 消费者端实现高可用的重要手段。2.2 系统架构一个典型的Kafka集群中包含若干Producer可以是web前端产生的Page View或者是服务器日志系统CPU、Memory等若干brokerKafka支持水平扩展一般broker数量越多集群吞吐率越高若干Consumer Group以及一个Zookeeper集群。Kafka通过Zookeeper管理集群配置选举leader以及在Consumer Group发生变化时进行rebalance。Producer使用push模式将消息发布到brokerConsumer使用pull模式从broker订阅并消费消息。2.3 生产者2.3.1 数据执行流程在 Kafka 中我们把产生消息的那一方称为生产者比如我们经常回去淘宝购物你打开淘宝的那一刻你的登陆信息登陆次数都会作为消息传输到 Kafka 后台当你浏览购物的时候你的浏览信息你的搜索指数你的购物爱好都会作为一个个消息传递给 Kafka 后台然后淘宝会根据你的爱好做智能推荐致使你的钱包从来都禁不住诱惑那么这些生产者产生的消息是怎么传到 Kafka 应用程序的呢发送过程是怎么样的呢尽管消息的产生非常简单但是消息的发送过程还是比较复杂的如图我们从创建一个ProducerRecord 对象开始ProducerRecord 是 Kafka 中的一个核心类它代表了一组 Kafka 需要发送的 key/value 键值对它由记录要发送到的主题名称Topic Name可选的分区号Partition Number以及可选的键值对构成。在发送 ProducerRecord 时我们需要将键值对对象由序列化器转换为字节数组这样它们才能够在网络上传输。然后消息到达了分区器。如果发送过程中指定了有效的分区号那么在发送记录时将使用该分区。如果发送过程中未指定分区则将使用key 的 hash 函数映射指定一个分区。如果发送的过程中既没有分区号也没有则将以循环的方式分配一个分区。选好分区后生产者就知道向哪个主题和分区发送数据了。ProducerRecord 还有关联的时间戳如果用户没有提供时间戳那么生产者将会在记录中使用当前的时间作为时间戳。Kafka 最终使用的时间戳取决于 topic 主题配置的时间戳类型。然后这条消息被存放在一个记录批次里这个批次里的所有消息会被发送到相同的主题和分区上。由一个独立的线程负责把它们发到 Kafka Broker 上。Kafka Broker 在收到消息时会返回一个响应如果写入成功会返回一个 RecordMetaData 对象它包含了主题和分区信息以及记录在分区里的偏移量上面两种的时间戳类型也会返回给用户。如果写入失败会返回一个错误。生产者在收到错误之后会尝试重新发送消息几次之后如果还是失败的话就返回错误消息。上面写的有点多总结一下流程创建对象主题、分区、key/value- 序列化数据 - 到达分区可自己指定也可以通过key hash- 放入批次相同主题和分区 - 独立线程发送 - 返回主题/分区/分区偏移量/时间戳。2.3.2 分区策略Kafka 对于数据的读写是以分区为粒度的分区可以分布在多个主机Broker中这样每个节点能够实现独立的数据写入和读取并且能够通过增加新的节点来增加 Kafka 集群的吞吐量通过分区部署在多个 Broker 来实现负载均衡的效果下面我们看看数据如何选择分区。方式1顺序轮询顺序分配消息是均匀的分配给每个 partition即每个分区存储一次消息见下图。轮训策略是 Kafka Producer 提供的默认策略如果你不使用指定的轮训策略的话Kafka 默认会使用顺序轮训策略的方式。方式2随机轮询本质上看随机策略也是力求将数据均匀地打散到各个分区但从实际表现来看它要逊于轮询策略所以如果追求数据的均匀分布还是使用轮询策略比较好。事实上随机策略是老版本生产者使用的分区策略在新版本中已经改为轮询了。方式3key hash这个策略也叫做 key-ordering 策略Kafka 中每条消息都会有自己的key一旦消息被定义了 Key那么你就可以保证同一个 Key 的所有消息都进入到相同的分区里面由于每个分区下的消息处理都是有顺序的故这个策略被称为按消息键保序策略如下图所示2.4 消费者2.4.1 消费者群组应用程序使用 KafkaConsumer 从 Kafka 中订阅主题并接收来自这些主题的消息然后再把他们保存起来。应用程序首先需要创建一个 KafkaConsumer 对象订阅主题并开始接受消息验证消息并保存结果。一段时间后生产者往主题写入的速度超过了应用程序验证数据的速度这时候该如何处理如果只使用单个消费者的话应用程序会跟不上消息生成的速度就像多个生产者像相同的主题写入消息一样这时候就需要多个消费者共同参与消费主题中的消息对消息进行分流处理。Kafka 消费者从属于消费者群组。一个群组中的消费者订阅的都是相同的主题每个消费者接收主题一部分分区的消息。下面是一个 Kafka 分区消费示意图。上图中的主题 T1 有四个分区分别是分区0、分区1、分区2、分区3我们创建一个消费者群组1消费者群组中只有一个消费者它订阅主题T1接收到 T1 中的全部消息。由于一个消费者处理四个生产者发送到分区的消息压力有些大需要帮手来帮忙分担任务于是就演变为下图这样一来消费者的消费能力就大大提高了但是在某些环境下比如用户产生消息特别多的时候生产者产生的消息仍旧让消费者吃不消那就继续增加消费者。如上图所示每个分区所产生的消息能够被每个消费者群组中的消费者消费如果向消费者群组中增加更多的消费者那么多余的消费者将会闲置如下图所示。向群组中增加消费者是横向伸缩消费能力的主要方式。总而言之我们可以通过增加消费组的消费者来进行水平扩展提升消费能力。这也是为什么建议创建主题时使用比较多的分区数这样可以在消费负载高的情况下增加消费者来提升性能。另外消费者的数量不应该比分区数多因为多出来的消费者是空闲的没有任何帮助。Kafka 一个很重要的特性就是只需写入一次消息可以支持任意多的应用读取这个消息。换句话说每个应用都可以读到全量的消息。为了使得每个应用都能读到全量消息应用需要有不同的消费组。对于上面的例子假如我们新增了一个新的消费组 G2而这个消费组有两个消费者那么就演变为下图这样。在这个场景中消费组 G1 和消费组 G2 都能收到 T1 主题的全量消息在逻辑意义上来说它们属于不同的应用。总结起来就是如果应用需要读取全量消息那么请为该应用设置一个消费组如果该应用消费能力不足那么可以考虑在这个消费组里增加消费者。2.4.2 消费者重平衡我们从上面的消费者演变图中可以知道这么一个过程最初是一个消费者订阅一个主题并消费其全部分区的消息后来有一个消费者加入群组随后又有更多的消费者加入群组而新加入的消费者实例分摊了最初消费者的部分消息这种把分区的所有权通过一个消费者转到其他消费者的行为称为重平衡英文名也叫做 Rebalance 。如下图所示。重平衡非常重要它为消费者群组带来了高可用性 和 伸缩性我们可以放心的添加消费者或移除消费者不过在正常情况下我们并不希望发生这样的行为。在重平衡期间消费者无法读取消息造成整个消费者组在重平衡的期间都不可用。另外当分区被重新分配给另一个消费者时消息当前的读取状态会丢失它有可能还需要去刷新缓存在它重新恢复状态之前会拖慢应用程序。消费者通过向组织协调者Kafka Broker发送心跳来维护自己是消费者组的一员并确认其拥有的分区。对于不同不的消费群体来说其组织协调者可以是不同的。只要消费者定期发送心跳就会认为消费者是存活的并处理其分区中的消息。当消费者检索记录或者提交它所消费的记录时就会发送心跳。如果过了一段时间 Kafka 停止发送心跳了会话Session就会过期组织协调者就会认为这个 Consumer 已经死亡就会触发一次重平衡。如果消费者宕机并且停止发送消息组织协调者会等待几秒钟确认它死亡了才会触发重平衡。在这段时间里死亡的消费者将不处理任何消息。在清理消费者时消费者将通知协调者它要离开群组组织协调者会触发一次重平衡尽量降低处理停顿。重平衡是一把双刃剑它为消费者群组带来高可用性和伸缩性的同时还有有一些明显的缺点(bug)而这些 bug 到现在社区还无法修改。重平衡的过程对消费者组有极大的影响。因为每次重平衡过程中都会导致万物静止参考 JVM 中的垃圾回收机制也就是 Stop The World STW。也就是说在重平衡期间消费者组中的消费者实例都会停止消费等待重平衡的完成而且重平衡这个过程很慢......2.5 特性分析这里才是内容的重点不仅需要知道Kafka的特性还需要知道支持这些特性的原因消息路由不支持Kafka在处理消息之前是不允许消费者过滤一个主题中的消息。一个订阅的消费者在没有异常情况下会接受一个分区中的所有消息。消息有序支持当消费消息时如果消费失败消息不会被放回所以整个消费过程都是有序进行消息时序不支持消息直接发送不会延迟发送或者指定消息的TTL。容错处理集群支持/消息不支持集群容错能力高因为是分布式部署但是消息容错处理弱因为消息消费失败需要程序员手动处理Kafka不支持消息重新进行消费。伸缩非常好通过扩充分区和消费者数量实现分区扩容并提升消费速度。持久化非常好数据存储在磁盘可以随时订阅消费消费完后数据仍然保留。消息回溯支持因为消息支持持久化就支持回溯可以理解是附带的功能。高吞吐非常好因为Kafka内部同一个主题包含多个分区所以实现分布式存储然后消费者数量可以扩充到和分区数量一致保证了Kafka的高吞吐。

关于恒美微站

恒美微站专注于为个体商户、工作室提供极简自助建站服务,让每个人都能轻松拥有专业网站。

快速链接

  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心

服务项目

  • 可视化建站
  • 拖拽编辑
  • 主题定制
  • SEO 优化
  • 网站托管

联系方式

  • 📍 地址:北京市朝阳区建国路 88 号
  • 📞 电话:400-888-8888
  • ✉️ 邮箱:info@hmyw.cn
  • 🕐 时间:周一至周日 9:00-18:00

© 2024 恒美微站 hmyw.cn 版权所有 | 京 ICP 备 12345678 号