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 offset1.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的所有信息。通常情况下这样的文件夹应该有50个编号从0到49。打开图中的任意一个文件夹会发现它就是一个正常的Kafkatopic日志文件目录里面至少有一个日志文件.log和两个索引文件.index 和.timeindex。该日志中保存的消息都是 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.mssession.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指定从最早的位移开始消费。注意这里最早的位移不一定就是0。2.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即不要关闭这些空闲连接。在zk的bin目录下启动客户端脚本查看节点信息:./zkCli.sh执行命令查看节点数据:ls /查看kafka集群配置信息:

相关新闻

数字孪生案例 | FPSO数字孪生模拟可视化:船体结构、生产工艺、系泊外输全流程

数字孪生案例 | FPSO数字孪生模拟可视化:船体结构、生产工艺、系泊外输全流程

一艘FPSO(浮式生产储卸油装置),日处理原油15到20万桶,载重10到30万吨,造价动辄数十亿元。它是海上油气开发的核心装备——把海底采出来的原油直接在海上完成处理、储存和外输,相当于一座"漂浮的炼油厂…

2026/7/30 17:00:18 阅读更多 →
springboot+redis+lua脚本进行接口限流,解决高并发计数不准确问题

springboot+redis+lua脚本进行接口限流,解决高并发计数不准确问题

为啥用redis呢(只是此处的使用原因):因为redis是一个内存数据库,效率高;redis支持事务;redis支持分布式,与系统无强关联,不管系统是单机还是分布式部署都支持。为啥用lua脚本呢&…

2026/7/30 16:59:18 阅读更多 →
如何快速提升编程体验:5个Maple Mono字体的终极技巧

如何快速提升编程体验:5个Maple Mono字体的终极技巧

如何快速提升编程体验:5个Maple Mono字体的终极技巧 【免费下载链接】maple-font Maple Mono: Open source monospace font with round corner, ligatures and Nerd-Font icons for IDE and terminal, fine-grained customization options. 带连字和控制台图标的圆角…

2026/7/30 16:59:18 阅读更多 →

最新新闻

企业数据安全防护体系构建与关键技术解析

企业数据安全防护体系构建与关键技术解析

1. 企业数据安全防护的现状与挑战现代企业运营过程中,数据资产已经成为最核心的竞争力之一。从客户信息、交易记录到商业机密、研发数据,这些数字资产的安全直接关系到企业的生存与发展。然而随着数字化转型的深入,数据安全威胁也呈现出新的特…

2026/7/30 17:05:20 阅读更多 →
SecureCRT加密密码遗忘解决方案:从原理到实践的完整恢复指南

SecureCRT加密密码遗忘解决方案:从原理到实践的完整恢复指南

1. 项目概述:当SecureCRT的加密连接密码成为“拦路虎”作为一名常年与服务器、网络设备打交道的运维工程师或开发者,SecureCRT这款终端仿真软件绝对是工具箱里的“老伙计”。它稳定、功能强大,尤其是其会话管理功能,能让我们把常用…

2026/7/30 17:05:20 阅读更多 →
驱动层透明加密实战:为文件数据穿上“隐形盔甲”

驱动层透明加密实战:为文件数据穿上“隐形盔甲”

1. 项目概述:当数据“裸奔”成为常态 在数字化办公成为标配的今天,我们每天都在与海量的非结构化数据打交道。一份即将提交的投标方案书、一组包含核心算法的源代码、一批记录着客户信息的影像扫描件,这些以文件、文档、影像形式存在的数字资…

2026/7/30 17:05:20 阅读更多 →
别再把小程序当线上价目表了!2026 年小程序的真正价值,90% 的老板都没摸透

别再把小程序当线上价目表了!2026 年小程序的真正价值,90% 的老板都没摸透

经常有老板说 “我做了小程序,根本没用,就是个摆着的电子菜单”,但深究下去就会发现,他们的小程序里就放了几张产品图、一个联系电话,既没有运营,也没有融入业务流程,自然没用。 很多人做不好小…

2026/7/30 17:05:20 阅读更多 →
智能机器人实验——机械臂画画

智能机器人实验——机械臂画画

实验器材Dobot机械臂笔小型摄像头电脑软件白纸项目描述创意来源:前段时间有个团队实现了通过手机拍照将照片传到机器中,机器实现画图。于是想着通过Dobot机械臂实现同样甚至更好的功能。在实验指导书上看到,Dobot机械臂能够简单的进行写字&am…

2026/7/30 17:05:20 阅读更多 →
苏州简易注销失败,企业该如何走一般注销流程?

苏州简易注销失败,企业该如何走一般注销流程?

一条路走不通,换条路也得走完——注销这件事,拖越久代价越大 很多老板觉得简易注销被驳回了,干脆就不管了。但真相是:不注销的后果,比注销麻烦一百倍。法人被限制高消费、银行账户被冻结、甚至影响子女读书——这些都不…

2026/7/30 17:04:20 阅读更多 →

日新闻

Windows驱动存储终极清理工具:DriverStoreExplorer完全指南

Windows驱动存储终极清理工具:DriverStoreExplorer完全指南

Windows驱动存储终极清理工具:DriverStoreExplorer完全指南 【免费下载链接】DriverStoreExplorer Driver Store Explorer 项目地址: https://gitcode.com/gh_mirrors/dr/DriverStoreExplorer 您是否曾因Windows系统盘空间不足而烦恼?是否遇到过设…

2026/7/30 0:00:13 阅读更多 →
如何3步掌握Video Download Helper:网页视频下载的完整实战指南

如何3步掌握Video Download Helper:网页视频下载的完整实战指南

如何3步掌握Video Download Helper:网页视频下载的完整实战指南 【免费下载链接】VideoDownloadHelper Chrome Extension to Help Download Video for Some Video Sites. 项目地址: https://gitcode.com/gh_mirrors/vi/VideoDownloadHelper 你是否曾经在浏览…

2026/7/30 0:00:13 阅读更多 →
“双减”后首个AI备课压力测试报告:覆盖32所中小学的176节AI辅助课,暴露4大隐性增负节点

“双减”后首个AI备课压力测试报告:覆盖32所中小学的176节AI辅助课,暴露4大隐性增负节点

更多请点击: https://intelliparadigm.com 第一章:AI 教师备课辅助 AI 教师备课辅助系统正逐步成为教育数字化转型的核心支撑工具,它并非替代教师,而是通过语义理解、知识图谱与多模态生成能力,将教师从重复性劳动中解…

2026/7/30 0:00:13 阅读更多 →

周新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/29 22:18:20 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/29 14:34:28 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/29 15:00:03 阅读更多 →

月新闻