惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

宝玉的分享
宝玉的分享
T
The Blog of Author Tim Ferriss
Y
Y Combinator Blog
Apple Machine Learning Research
Apple Machine Learning Research
Last Week in AI
Last Week in AI
Recorded Future
Recorded Future
博客园 - 司徒正美
V
Vulnerabilities – Threatpost
月光博客
月光博客
C
CXSECURITY Database RSS Feed - CXSecurity.com
D
Darknet – Hacking Tools, Hacker News & Cyber Security
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
Microsoft Azure Blog
Microsoft Azure Blog
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
W
WeLiveSecurity
Jina AI
Jina AI
Exploit-DB.com RSS Feed
Exploit-DB.com RSS Feed
Hacker News: Ask HN
Hacker News: Ask HN
S
Security Affairs
V
Visual Studio Blog
Schneier on Security
Schneier on Security
T
Tailwind CSS Blog
Martin Fowler
Martin Fowler
V2EX - 技术
V2EX - 技术
博客园 - Franky
S
Secure Thoughts
Blog — PlanetScale
Blog — PlanetScale
G
GRAHAM CLULEY
D
DataBreaches.Net
O
OpenAI News
Forbes - Security
Forbes - Security
云风的 BLOG
云风的 BLOG
Google Online Security Blog
Google Online Security Blog
博客园 - 三生石上(FineUI控件)
T
Tor Project blog
T
Tenable Blog
Latest news
Latest news
N
News and Events Feed by Topic
博客园_首页
Simon Willison's Weblog
Simon Willison's Weblog
C
Cybersecurity and Infrastructure Security Agency CISA
美团技术团队
T
The Exploit Database - CXSecurity.com
K
Kaspersky official blog
B
Blog
阮一峰的网络日志
阮一峰的网络日志
T
Threat Research - Cisco Blogs
SecWiki News
SecWiki News
PCI Perspectives
PCI Perspectives
GbyAI
GbyAI

博客园_首页

Linux实操--组管理、权限管理和定时任务 Java + EasyExcel 实现单个接口导出多个Excel Mem0 源码解析系列(二):提示词工程的深度剖析 Openclaw TaskFlow究竟是什么?和普通Skill技能有什么区别 博文阅读密码验证 - 博客园 嘉立创开源:应该是全网MicroPython教程最多的开发板 Hermes Agent 集成实践:从协议到生产 2026年AI编程工具横评:Cursor、Codex、Claude Code、Zed、Windsurf Java程序员必看的RAG入门教程 2026 AI效率神器:Superpowers + Claude Code 保姆级教程 本地大模型部署全攻略:从 0 到 1 玩转 Ollama 【从0到1构建一个ClaudeAgent】内存管理-上下文压缩 .NET 高级开发 | 设计、实现一个事件总线框架 电子小白入门之NE555 3. WorkBuddy:隐藏玩法,一键召唤专家,让 AI 以"专家身份"给你干活 和AI一起搞事情#3:Claude Teammate 游戏开发翻车实录 【OpenClaw】通过 Nanobot 源码学习架构---(7)Memory C# .NET 周刊|2026年3月3期 我在 Debian 11 上把 K8s 单机搭起来了,过程没你想的那么顺(/opt 目录版) 深度学习进阶(七)Data-efficient Image Transformer CLI+Skill搭建浏览器AI自动化框架,告别一切重复枯燥任务 告别Token账单无底洞:OpenClaw本地部署,重塑企业数据主权的唯一解 FastAPI+Vue:文件分片上传+秒传+断点续传,这坑我帮你踩平了! SBTI 爆火后,我做了个程序员版的 CBTI。。已开源 + 附开发过程 多模态检索开始进入工程期:用 Sentence Transformers 搭建可落地的 Multimodal RAG 100多行代码实现一个最简单的Agent(用ReAct) Claude Code 通关手册(八):推荐 5 个 Hooks,代码质量提升 3 倍 老板:“有人截图了!”。安全部门:“收到,马上查暗水印!” - why技术 技术之外,皆是人间 C#/.NET/.NET Core技术前沿周刊 | 第 69 期(2026年4.01-4.12) Snack JSONPath 项目架构分析 Claude Code Buddy 小析:一个非核心功能,如何体现产品的细节完成度 AI新时代下的图床管理方案-Cloudflare图床+MCP+Skills方案指南 化繁为简:顺丰速运App如何通过 HarmonyOS SDK实现专业级空间测量 从零实现富文本编辑器#13-React非编辑节点的内容渲染 AI开发-python-langchain框架(3-23-OpenAI Functions风格Tool Calling智能助手) .NET + AI 进阶实战:基于类的技能开发 - 打造可治理的 Agent 能力模块 【从0到1构建一个ClaudeAgent】规划与协调-技能 上周热点回顾(4.6-4.12) 电子小白的工具三件套:面包板、杜邦线、万能板 单表五亿数据的查询优化 | Mysql、StarRocks 2. WorkBuddy:从“我是谁”到“帮我干活” C# 如何减少代码运行时间:7 个实战技巧 基于HelixToolkit.SharpDX 渲染3D模型 - 笺上知微 从零开始的双臂具身VLA起源及现阶段发展综述 - SkyXZ 记对 xonsh shell 的使用, 脚本编写, 迁移及调优 - pluvium27 受够了Vibe Coding的失控?换个起点,让AI事半功倍 从开始配置漏洞环境到漏洞复现流程 - 難しい 关于10年工作经验的程序员对OpenClaw的实战经验分享以及看法 - 虚无境 Any metadata 的内存布局 C# .NET 周刊|2026年3月2期 - InCerry 我帮你测过了,测试圈排名第二的 Skill 依然很牛逼 Skill Discovery | 无监督技能发现的经典工作总结 - MoonOut PbootCMS 网站内容数量多导致访问慢?这些实用优化方案帮你提速! - 家兴网络技术工作室 上下文工程是什么?过时了么?一文讲明白! - 一枫说码 网站漏洞怎么发现并修复?一篇实用指南(附完整流程) - 家兴网络技术工作室 开了 TUN 模式还是直连?90% 的人都踩过这个坑 Github日报|2026年04月12日 - AI一族 AScript扩展多种脚本语言 - rockey627 AI 学习笔记:Agent 的记忆机制 你能被装进一个文件里吗?——7 万人把同事"蒸馏"成了 AI - 我没有三颗心脏 Claude Code 通关手册(七):给 AI 装上技能包——Skills 完全指南 - 暮色之狐 在浏览器中快速编辑代码:VSCode Web 集成实践 - Newbe36524 蒸馏自己 skill?基于 Deepseek 的蒸馏器,丐版蒸馏方式,简单便捷 - To_Carpe_Diem Spring AI Aliababa和AgentScope,哪个更好? - 苏三说技术 Etsy 把 1000 个 MySQL 分片迁进 Vitess:425TB 数据背后的真正问题不是性能,而是运维规模 MicroPython LVGL基础知识和概念:底层渲染与性能优化 - FreakStudio 数据库草图算法 Python 潮流周刊#146:CPython 引入 Rust 的进展 - 豌豆花下猫 最小生成树 - mofei1116 红日靶场七:从外网入口、容器逃逸到 AD 接管的完整利用链复盘 - YouDiscovered1t 分享四款开源且实用的 Kafka 管理工具 - 追逐时光者 vLLM 权重加载机制全解析:从挑战到理想架构 LCT 学习笔记 - ACehomoxue Avalonia UI 12.0.0 正式发布:架构演进和性能飞跃 - 张善友 当 AI Agent 把调用链拉长,延迟开始成为一门生意 conhost.exe 无法显示 U+2717 - 145a 太秀了,我把自己蒸馏成了 Skill!已开源 - 程序员鱼皮 ASP.NET Core 内存缓存实战:一篇搞懂该怎么配、怎么避坑 基于 Ghostty 带有分割标签页和为 Claude 编程设计的通知终端 - BugShare AI 焊死入口:教育的“操作系统级”重塑 - 郝hai 初级Java开发工程师使用sql脚本编写代码的过程是简单而且不糊涂 - CoderOilStation Claude Code通关手册(六):MCP协议完全指南 - 暮色之狐 边框灯光环绕动画特效实现指南 - Newbe36524 开源:子木蒸馏版的 SEO 审计工具 seo-audit-skill v1.0 我所理解的Python元模型 【从0到1构建一个ClaudeAgent】规划与协调-TodoWrite - 程序员Seven Claude 和 Codex 在审计 Skill 上性能差异探究 - ACai_sec AScript如何实现中文脚本引擎 - rockey627 【渗透测试】HTB Season10 Garfield 全过程wp - dynasty_chenzi Android 开发者为什么必须掌握 AI 能力?端侧视角下的技术变革 树状数组正确性证明 - AC-wyr 你的 AI 焦虑,可能比 AI 本身更危险——ATM 机没有消灭银行柜员,但恐慌消灭了你的判断力 - 我没有三颗心脏 一个拉胯的分库分表方案有多绝望?整个部门都在救火! - 冰河团队 动态规划入门必学之走方格问题 - Ofnoname PostgREST 与 PostgreSQL 角色权限配置全解析(生产级实践) - SheepDog1998 使用 UEFI 图形输出协议 GOP 在屏幕上显示图像的方法 - 阿源- Claude Code通关手册(五):组建你的AI专家团队,子代理系统 - 暮色之狐 一个程序员到架构师的催婚路之感悟(整整10年后的催婚相亲感悟) - MisterLip 用 Agent Skill 自动生成工作周报 - 赵康
从零学习Kafka:消费者组重平衡
Jackeyzhe · 2026-05-18 · via 博客园_首页

本文我们一起来学习消费者组重平衡相关的知识。

写在前面

先解决前面留下的问题:如果生产者已经发送了大量消息,但在最后提交之前突然宕机,事务协调器会如何处理这个未完成的事务呢?

答案是自动终止,事务协调器如果在一段时间内没有收到生产者的任何消息或者提交事务的请求,会利用 __transaction_state 中记录的信息向所有涉及到的分区发送 Abort 命令,写入的消息也会被标记为废弃。这个时长是由 transaction.timeout.ms 参数控制的,默认是 1 分钟(值为 60000)。

什么是消费者组

回答完遗留的问题,我们进入今天的正题,第一个问题是什么是消费者组(Consumer Group)?

这是 Kafka 中比较有亮点的设计。简单来说,Consumer Group 就是 Kafka 提供的可扩展且具有容错性的消费机制。消费者组内部包含多个消费者实例,它们共享一个 Group ID。组内的所有消费者实例共同订阅一个或多个 Topic 的所有分区。每个分区只能由同一个消费者组的一个消费者实例来消费。

核心特性

消费者组有以下特性:

  • 负载均衡:Kafka 会把订阅的分区均衡的分配给组内的所有消费者,这样就可以实现并行处理了。理想情况下,消费者数量应该等于消费组订阅主题的分区总数。
  • 点对点模式:还记得之前我们讲过 Kafka 既支持点对点模型,又支持发布/订阅模型吗?实际就是通过 Consumer Group 来支持的。如果所有消费者都在同一个组内,每条消息都只会由其中一个消费者处理,类似于传统的消息队列。这就是点对点模型的实现。
  • 发布/订阅模型:如果多个不同的消费者组订阅同一个 Topic,那么每个组都会收到一份完整的消息,互相之间不会干扰。这种就是 发布/订阅模型。
  • 故障转移:如果消费组内某个实例宕机了,Kafka 会自动检测到,并触发重平衡。将其负责的分区分配给正常的实例,以此保证数据不丢失。

消费者组重平衡流程

触发条件

聊完消费者组的定义和特性之后,我们再来看一下它最核心,也是最让人头疼的机制——重平衡。重平衡本质上是一种自动调度的机制,它确保 Topic 的所有分区都有对应的消费者来消费,并且分配的相对均衡。

触发重平衡的条件有三个:

  1. 组员数量发生变动:有新消费者加入组,或者某个实例崩溃。

  2. 订阅主题数变动:消费者组可以使用正则表达式订阅主题,当创建了新的满足条件的主题时,就会触发。

  3. 订阅主题的分区数发生变化:订阅的主题扩展分区。

消费者组的状态

Kafka 为消费者组定义了 9 种状态,它们的含义如下:

状态 含义
UNKNOWN 未知状态,通常是客户端与 Broker 的版本不兼容,或者发生了无法识别的异常。
PREPARING_REBALANCE 重平衡准备中,协调器收到了加入申请或者心跳超时,准备重新分配。
COMPLETING_REBALANCE 重平衡同步中,所有成员都已经加入组,Leader 正在计算协调方案。
STABLE 稳定状态,重平衡完成。
DEAD 消费者组注销,元信息在协调者端已被移除。
EMPTY 组内没有任何成员,通常是刚创建或者所有的消费者都正常关闭,元信息依然保留在协调者端。
ASSIGNING 分配中,这是 2.4+ 版本新引入的,Leader 计算好了一部分新的分配计划,正准备下发。
RECONCILING 协调中/一致性同步中,也是 2.4+ 版本新引入的,成员正在释放不再属于自己的分区,准备接手新的分区。
NOT_READY 未就绪,通常是组刚刚启动,或者协调者端在进行迁移。

状态流转流程如下图

KafkaGroupState

接着我们来看一下最经典的重平衡过程。

  1. 正常情况下消费者组的状态是 STABLE(也可能是 Empty),协调者与消费者组中的组员之间维护有心跳消息。
  2. 当有一个新的成员要加入时,会给协调者发送一个 JoinGroup 请求。协调者收到请求后,会通过心跳消息通知其他消费者。所有消费者此时会停止消费,重新向协调者发送 JoinGroup 请求,消费者组进入 PREPARING_REBALANCE 状态。
  3. 当所有成员都到齐之后,协调者会从中选出一个作为 Leader,然后把所有的消费者信息通过 JoinGroup 的响应发送给 Leader。此时消费者组为 COMPLETING_REBALANCE 状态。
  4. Leader 计算好分配方案后,会通过 SyncGroup 请求,将其发送给协调者。其他的成员也会发送 SyncGroup 请求,这是为了方便协调者将分配方案包装进 SyncGroup 的响应中返回给所有的消费者。
  5. 消费者拿到新的任务之后,就开始继续工作了。消费者组的状态恢复成 STABLE。

这就是一次完整的重平衡流程,这里有一个问题是:在这整个过程中,消费者都是不处理消息的,也就是我们常说 Stop-The-World 问题。如果你有几百个 Consumer 实例,那么一次 Rebalance 可能需要几个小时,这简直令人崩溃。

基于这种问题,Kafka 在 2.4 版本推出了增量协作重平衡机制。这种机制下,重平衡的过程不再要求所有的消费者都放下手中的工作,而是只处理那些需要变动的分区,这样就极大的提升了稳定性。只需要在大于 2.4 版本的客户端中将 partition.assignment.strategy 设置为 CooperativeStickyAssignor 即可。

新版本重平衡流程如下:

  1. 协调者检测到变动时,开启第一轮协商,此时状态由 STABLE 变为 PREPARING_REBALANCE。

  2. 所有成员到齐后,协调者把消费者的信息发送给 Leader,消费者组状态变为 COMPLETING_REBALANCE。

  3. Leader 算出哪些分区需要释放,并通知给相关消费者,此时消费者组状态变为 RECONCILING。

  4. 释放完成后,消费者组状态回到 STABLE,此时存在一些游离分区需要认领。

  5. 协调器接着开启新一轮的协商,通过相同的报道步骤,状态从 STABLE 变为 PREPARING_REBALANCE 再变为 COMPLETING_REBALANCE。

  6. 这一阶段主要目的是为游离分区进行重新分配,此时状态变成 ASSIGNING。

  7. 当所有游离分区都有消费者认领后,再次回到稳定的 STABLE 状态。

Kafka 通过多次小的调整,来避免整个集群长时间停止工作,以此来减少重平衡对于整体集群的影响。这一进化是不是有点像 JVM 的 GC 从传统垃圾回收器进化到 G1 和 ZGC。

如何避免重平衡

虽然新版本的重平衡机制有了很大的进步,但还是会对系统性能造成一定的影响。那如何才能避免重平衡呢?

首先,完全消除重平衡是不可能的。我们要做的就是消除掉非预期的重平衡。什么是非预期的呢?你可以理解为是由于配置不当或者系统抖动引起的重平衡。

我们分别从参数调优、代码健壮性和架构设计三个层面来看一下如何调整。

参数调优

首先是参数调优,非预期重平衡触发最常见的两个原因,一个是心跳超时,另一个是逻辑处理超时。

为了避免因为网络抖动导致误判心跳超时,我们可以适当调大 session.timeout.ms,这个参数决定了 Consumer 存活性的时间间隔,除了这个参数,还需要调整 heartbeat.interval.ms,这个是用来控制发送心跳消息的频率的。发送的越频繁,协调者越能更快响应 Consumer 掉线并开启重平衡,但随之而来的问题是消耗的资源也越多。通常可以把它设置为 session.timeout.ms 的三分之一。

逻辑处理超时的参数主要是 max.poll.interval.ms,它用来控制两次 poll 之间的间隔,如果你的业务逻辑复杂,需要处理时间比较长,那么就需要调大这个参数。例如你在业务代码中访问了第三方存储,整个过程需要 5 分钟,那么这个参数可以设置为 6 分钟。除了调大 max.poll.interval.ms 之外,我们也可以调整 max.poll.records,它是用来控制每次 poll 的消息条数,通过减少消息条数,能够缩短 poll 一次的逻辑处理时间。

代码健壮性

介绍完了参数调优之后,我们再来看一下代码层面有哪些需要调整或者注意的地方。首先是最基本的 try-catch 防护,我们应该确保所有的异常都在消费逻辑中处理掉,一旦消费者因为没捕获异常而崩溃,那么必然会触发重平衡。

其次就是优雅关闭 Consumer,在停止时手动调用 consumer.close(),这样会给协调者发送 LeaveGroup 请求,协调者收到请求后可以立即开启重平衡,缩短“空窗期”。

架构设计

在架构层面,除了使用我们前面提到的增量协作重平衡协议之外,还可以设置 group.instance.id,这是为每个消费者实例设置一个固定的 ID,这样在实例重启时,只要在 session.timeout.ms 时间内回来,协调者都会认出它,不会触发重平衡。

总结

本文我们先了解了什么是消费者组,一句话概括就是它是 Kafka 提供的可扩展且具有容错性的消费机制。接着又聊了重平衡机制,包括消费者组的状态以及重平衡的整个流程。最后我们介绍了如何避免非预期的重平衡,这能帮助我们提升 Kafka 集群的稳定性。