当前位置: 首页 > news >正文

7 kafka在zk中的目录

1. 为什么需要consumer group 好处是什么?

1.实际上,consumer group是用于实现高伸缩性、高容错性的consumer机制。2.组内多个 consumer实例可以同时读取Kafka消息,而且一旦有某个 consumer“挂”了,consumer group会立即将已崩溃 consumer负责的分区转交给其他 consumer来负责,从而保证整个group可以继续工作,不会丢失数据——这个过程被称为重平衡(rebalance)。

2 消费组和消息的顺序性关系

另外由于 Kafka目前只提供单个分区内的消息顺序,而不会维护全局的消息顺序,因此如果用户要实现 topic 全局的消息读取顺序,就只能通过让每个 consumer group 下只包含一个consumer实例的方式来间接实现。

3 consumer offset

1.consumer端的 offset,与分区日志中的 offset是不同的含义。2.每个 consumer 实例都会为它消费的分区维护属于自己的位置信息来记录当前消费了多少条消息。

很多消息引擎都把消费端的 offset 保存在服务器端(broker),这样做的好处当然是实现简单,但会有以下3个方面的问题。

1.broker从此变成了有状态的,增加了同步成本,影响伸缩性2.需要引入应答机制(acknowledgement)来确认消费成功3.由于要保存许多 consumer 的 offset,故必然引入复杂的数据结构,从而造成不必要的资源浪费。

Kafka则选择了不同的方式:让 consumer group保存 offset,那么只需要简单地保存一个长整型数据就可以了,同时 Kafka consumer 还引入了检查点机制(checkpointing)定期对offset进行持久化,从而简化了应答机制的实现。

从下图中我们可以看到当前Kafka consumer在内部使用一个 map来保存其订阅topic所属分区的offset。

4 offset提交

consumer客户端需要定期地向Kafka集群汇报自己消费数据的进度,这一过程被称为位移提交(offset commit)。

新版本和旧版本 consumer提交位移的方式截然不同:旧版本 consumer会定期将位移信息提交到ZooKeeper下的固定节点上

把位移提交到ZooKeeper的做法并不合适。ZooKeeper本质上只是一个协调服务组件,它并不适合作为位移信息的存储组件,毕竟频繁高并发的读/写操作并不是 ZooKeeper擅长的事情。

新版本consumer把位移提交到 Kafka 的一个内部 topic(__consumer_offsets)上,通常不能直接操作该topic就可以了,特别是注意不要擅自删除或搬移该topic的日志文件。

5 _consumer_offsets

_consumer_offsets是Kafka自行创建的,因此用户不可擅自删除该 topic的所有信息。

  1. 通常情况下,这样的文件夹应该有50个,编号从0到49。
  2. 打开图中的任意一个文件夹,会发现它就是一个正常的Kafkatopic日志文件目录,里面至少有一个日志文件(.log)和两个索引文件(.index 和.timeindex)。
  3. 该日志中保存的消息都是 Kafka 集群上consumer(特别是 consumer group)的位移信息罢了。

_consumer_offsets的每条消息格式大致如图

你可以把它想象成一个KV格式的消息,key就是一个三元组:group.id + topic + 分区号,而value就是offset的值。每当更新同一个 key的最新offset值时,该topic就会写入一条含有最新 offset的消息,同时 Kafka会定期地对该 topic执行压实操作(compact),即为每个消息key 只保存含有最新offset的消息。这样既避免了对分区日志消息的修改,也控制住了_consumer_offsets topic总体的日志容量,同时还能实时反映最新的消费进度。

考虑到一个Kafka生产环境中可能有很多consumer或consumer group,如果这些consumer同时提交位移,则必将加重__consumer_offsets的写入负载,因此社区特意为该topic创建了50个分区,并且对每个group.id做哈希求模运算,从而将负载分散到不同的__consumer_offsets分区上。这就是说,每个consumer group保存的offset都有极大的概率分别出现在该topic的不同分区上。

6 消费者组重平衡

如果使用的是 standalone consumer,则压根就没有rebalance的概念,即rebalance只对consumer group有效。

何为 rebalance?它本质上是一种协议,规定了一个 consumer group下所有 consumer如何达成一致来分配订阅 topic的所有分区。举个例子,假设我们有一个 consumer group,它有20个 consumer实例。该 group订阅了一个具有100个分区的 topic。那么正常情况下,consumer group平均会为每个consumer分配5个分区,即每个 consumer负责读取5个分区的数据。这个分配过程就被称作rebalance。

7 consumer主要参数

7.1 session.timeout.ms

session.timeout.ms是consumer group检测组内成员发送崩溃的时间。

假设你设置该参数为5分钟,那么当某个group成员突然崩溃了(比如被kill-9或宕机),管理 group的 Kafka 组件(即消费者组协调者,也称 group coordinator,有可能需要 5 分钟才能感知到这个崩溃。

显然我们想要缩短这个时间,让coordinator 能够更快地检测到 consumer 失败。遗憾的是,这个参数还有另外一重含义:consumer消息处理逻辑的最大时间——倘若consumer两次poll之间的间隔超过了该参数所设置的阈值,那么coordinator 就会认为这个 consumer 已经追不上组内其他成员的消费进度了,因此会将该consumer实例“踢出”组,该consumer负责的分区也会被分配给其他consumer。

在最好的情况下,这会导致不必要的rebalance,因为consumer需要重新加入group。

更糟的是,对于那些在被踢出group后处理的消息,consumer都无法提交位移——这就意味着这些消息在rebalance之后会被重新消费一遍。

如果一条消息或一组消息总是需要花费很长的时间处理,那么consumer甚至无法执行任何消费,除非用户重新调整参数。

鉴于以上的“窘境”,Kafka社区于0.10.1.0版本对该参数的含义进行了拆分。在该版本及以后的版本中,session.timeout.ms 参数被明确为“coordinator 检测失败的时间”。因此在实际使用中,用户可以为该参数设置一个比较小的值让 coordinator能够更快地检测 consumer崩溃的情况,从而更快地开启 rebalance,避免造成更大的消费滞后(consumer lag)。目前该参数的默认值是10秒。

7.2 max.poll.interval.ms

如前所述,session.timeout.ms 中“consumer 处理逻辑最大时间”的含义被剥离出来了,Kafka为这部分含义单独开放了一个参数——max.poll.interval.ms。

通过将该参数设置成实际的逻辑处理时间再结合较低的session.timeout.ms 参数值,consumer group既实现了快速的consumer崩溃检测,也保证了复杂的事件处理逻辑不会造成不必要的rebalance。

7.3 auto.offset.reset

指定了无位移信息或位移越界(即 consumer 要消费的消息的位移不在当前消息日志的合理区间范围)时 Kafka的应对策略。

特别要注意这里的无位移信息或位移越界,只有满足这两个条件中的任何一个时该参数才有效果。

举例说明:
假设你首次运行一个consumer group并且指定从头消费。显然该group会从头消费所有数据,因为此时该 group 还没有任何位移信息。一旦该 group 成功提交位移后,你重启了 group,依然指定从头消费。此时你会发现该 group并不会真的从头消费——因为Kafka已经保存了该group的位移信息,因此它会无视auto.offset.reset的设置。

该参数有如下3个可能的取值

1。earliest:指定从最早的位移开始消费。注意这里最早的位移不一定就是02.latest:指定从最新处位移开始消费3.none:指定如果未发现位移信息或位移越界,则抛出异常。在实际使用过程中几乎从未见过将该参数设置为none的用法,因此该值在真实业务场景中使用甚少。

7.4 enable.auto.commit

该参数指定 consumer是否自动提交位移。若设置为 true,则 consumer在后台自动提交位移;否则,用户需要手动提交位移。

7.5 fetch.max.bytes

指定了 consumer 端单次获取数据的最大字节数。若实际业务消息很大,则必须要设置该参数为一个较大的值,否则consumer将无法消费这些消息。

7.6 max.poll.records

该参数控制单次 poll调用返回的最大消息数。比较极端的做法是设置该参数为1,那么每次 poll只会返回1条消息。如果用户发现 consumer端的瓶颈在 poll速度太慢,可以适当地增加该参数的值。如果用户的消息处理逻辑很轻量,默认的500条消息通常不能满足实际的消息处理速度。

7.7 heartbeat.interval.ms

要搞清楚consumergroup的其他成员如何得知要开启新一轮rebalance——当coordinator决定开启新一轮rebalance时,它会将这个决定以REBALANCE_IN_PROGRESS异常的形式“塞进”consumer心跳请求的response中,这样其他成员拿到response后才能知道它需要重新加入group。显然这个过程越快越好,而heartbeat.interval.ms就是用来做这件事情的。

比较推荐的做法是设置一个比较低的值,让 group 下的其他 consumer成员能够更快地感知新一轮rebalance开启了。注意,该值必须小于session.timeout.ms!这很容易理解,毕竟如果consumer在session.timeout.ms这段时间内都不发送心跳,coordinator就会认为它已经dead,因此也就没有必要让它知晓coordinator的决定了。

7.8 connections.max.idle.ms

经常有用户抱怨在生产环境下周期性地观测到请求平均处理时间在飙升,这很有可能是因为 Kafka会定期地关闭空闲Socket连接导致下次consumer处理请求时需要重新创建连向broker的Socket连接。当前默认值是9分钟,如果用户实际环境中不在乎这些Socket资源开销,比较推荐设置该参数值为-1,即不要关闭这些空闲连接。

  1. 在zk的bin目录下,启动客户端脚本,查看节点信息:
    ./zkCli.sh

    执行命令查看节点数据:ls /

    查看kafka集群配置信息:
http://www.jsqmd.com/news/1294969/

相关文章:

  • 2026 年 7 月最新昆明防水补漏实测测评|卫生间、屋顶、外墙漏水维修靠谱商家横向对比 - 吉林同城获客
  • 2026犬步态分析设备推荐:犬足底压力系统厂家及选型指南 - 品牌深度评测
  • 企业级RAG应用:GPT-5.6长上下文检索与推理闭环实测
  • 收藏!AI Coding时代来临,小白程序员如何抓住高薪Agent工程师红利?
  • 单片机毕设选题推荐:基于嵌入式的点滴数据采集与电机调速系统设计 基于 STM32 的液体滴落检测与自动控制系统实现(013801)
  • 二维矩阵中的高效二分查找实现与优化
  • 2026年发物流用什么平台便宜?实测对比5家后我只推荐这个 - 快递物流资讯
  • 2026年苏州企业做AI搜索优化选哪家?牛橙网络科技实战解析 - 小子虾仁很开心
  • 2026年实力之选:专业管道电预热设备工程公司——新疆泓浩机电设备有限公司 - 企业推荐官【官方】
  • 暗黑2存档编辑器d2s-editor:三大革命重塑角色定制体验
  • 中国大模型全球份额碾压式领先:数据之外,我们该冷静看什么?
  • Windows 11终极清理指南:3分钟让你的系统重获新生
  • 2026年3款VIVO短视频总结哪个好,亲测对比后告诉你该怎么选
  • Modern-Screenshot深度解析:现代Web截图解决方案的架构设计与最佳实践
  • GetQzonehistory:5分钟快速恢复QQ空间历史数据的完整指南
  • 深圳南山区管道疏通避坑指南2026年7月找本地靠谱师傅 - 余生黄金回收
  • 大连企业必看:2026政府采购标书代做常见问题与解析 - 品牌优选官
  • AI搜索如何重构法律咨询流程?揭秘2024年律所已悄悄部署的5个智能检索实战案例
  • 计算机毕业设计之基于SpringBoot+Vue实现前后端分离商城管理系统的设计与实现
  • 百度网盘提取码智能获取工具:5秒快速破解加密资源的完整指南
  • 2026从一人直播到搭建工作室,我的六年踩坑与转型实录 - 彭拜新闻(测评)
  • 2026年椭圆机品牌推荐:五款热门横评 - 科技焦点
  • 界面控件DevExpress VCL v26.1新版亮点 - 支持Fluent UI
  • 2026年成都市成华区水电维修选维小达 电路维修、水管漏水抢修、管道疏通、马桶维修、暖气维修一站式服务 - 一点传媒
  • CUPS打印系统终极指南:7个关键场景下的深度配置与实战技巧
  • GHelper完全指南:告别臃肿,用轻量级工具掌控你的华硕笔记本
  • UE5集成OpenCV完整指南:从环境配置到实时图像处理
  • 3个技术方案解决桌面互动痛点:BongoCat如何用跨平台桌宠提升工作效率
  • AI音乐和弦进行效率革命:从手动编排到实时生成,97%专业制作人已在用的3个开源工具链
  • 便利店AI推荐优化:提升地图搜索排名的智能策略