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

推荐订阅源

N
Netflix TechBlog - Medium
Recorded Future
Recorded Future
云风的 BLOG
云风的 BLOG
Threat Intelligence Blog | Flashpoint
Threat Intelligence Blog | Flashpoint
T
Troy Hunt's Blog
Security Latest
Security Latest
Scott Helme
Scott Helme
Project Zero
Project Zero
C
Cybersecurity and Infrastructure Security Agency CISA
Google DeepMind News
Google DeepMind News
NISL@THU
NISL@THU
Latest news
Latest news
L
Lohrmann on Cybersecurity
P
Palo Alto Networks Blog
V
Vulnerabilities – Threatpost
Application and Cybersecurity Blog
Application and Cybersecurity Blog
AWS News Blog
AWS News Blog
S
Schneier on Security
TaoSecurity Blog
TaoSecurity Blog
C
Cyber Attacks, Cyber Crime and Cyber Security
N
News and Events Feed by Topic
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
H
Hacker News: Front Page
D
DataBreaches.Net
大猫的无限游戏
大猫的无限游戏
W
WeLiveSecurity
Cyberwarzone
Cyberwarzone
T
Threat Research - Cisco Blogs
The Register - Security
The Register - Security
L
LINUX DO - 热门话题
阮一峰的网络日志
阮一峰的网络日志
T
Tailwind CSS Blog
B
Blog
D
Docker
Apple Machine Learning Research
Apple Machine Learning Research
P
Privacy & Cybersecurity Law Blog
Engineering at Meta
Engineering at Meta
A
About on SuperTechFans
Webroot Blog
Webroot Blog
Martin Fowler
Martin Fowler
The Cloudflare Blog
Exploit-DB.com RSS Feed
Exploit-DB.com RSS Feed
S
SegmentFault 最新的问题
Vercel News
Vercel News
WordPress大学
WordPress大学
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
Recent Announcements
Recent Announcements
F
Fortinet All Blogs
Microsoft Azure Blog
Microsoft Azure Blog

博客园_首页

Plist 二进制格式 Milvus 和 PGVector,哪个更好? OpenClaw 已过时?在 VS Code 中运行 Hermes Agent! 第30篇文章:一个大三计科生的自白 Manim如何在数学公式中完美显示中文? Docker 部署 RocketMQ 5 并发编程核心概念辨析 C#事务处理最佳实践:别再让“主表存了、明细丢了”的破事发生 CLI 是什么?为什么大厂突然集体卷命令行? 【从0到1构建一个ClaudeAgent】协作-自主Agent UIImageView 设置图片不生效的原因排查 最小二乘问题详解20:无先验约束下的增量式SFM自由网平差 痞子衡嵌入式:大话双核i.MXRT1180之XIP应用里借助MU实现可靠Flash IAP的方法 AI Chat 封装, SemanticKerne.AiProvider.Unified 已发布 Windows下右键编辑js文件无法打开记事本——在注册表中使用环境变量 在后台服务中使用 Scoped 服务,为什么总是报错? H200 安装驱动并使用sglang启动模型 wireshark 抓包Trap上报告警内容 我用 AI 辅助开发了一系列小工具(2):图片压缩工具 [A Primer On MC and CC] 2.1 Memory Consistency 1 - 指令重排序和 SC 模型 Oracle数据库SCN推进技术详解与实践指南 玩转控件:封装个带图片的Label控件 Claude Code 4.7 真正该升级的不是模型,而是你的工作流 前端小白一句话,AI 帮我做了个颜值拉满的桌面媒体播放器。当代码不再是门槛,一句话编程就是现实。 5. WorkBuddy: 小龙虾的灵魂三件套,让你的小龙虾不只是工具 SQLite 分片方案实战:三种分片策略的深度对比 告别简陋 UI!一款基于 Fluent Design 和基于 WinUI 的开源免费、现代化的 Avalonia UI 控件库 关于二进制排列组合枚举的总结 AI开发-python-LangGraph框架(3-27-LangGraph从零实现大模型智能决策工作流) ElasticSearch主分片和副本分片概念详解 【002】HTTPS 粗解:证书、TLS 握手与对后端配置的影响 Hermes Agent 一周暴涨五万 Star,但我劝你别急着追 明明连接的是Redis的DB0,为什么能查到DB3的数据? 【从0到1构建一个ClaudeAgent】协作-Agent团队 熟悉电子元器件之后,电子小白下一步该怎么走? MAF快速入门(23)通过C#类定义Skills .NET 高级开发 | 手写一个对象映射框架 FastAPI数据库ORM怎么选?我肝了三个Demo后,终于不再纠结了 mysqldump 参数拾遗:在遗忘与铭记之间 C# .NET 周刊|2026年3月5期 Claude code入门 - 陈彦斌 一文学习入门 ThingsBoard 开源物联网平台 GitHub 热门项目 | 2026年04月16日 如何为GIT设置全局勾子,为每次提交追加信息 Number.isFinite和isFinite与isNaN()和Number.isNaN的区别 PortSwigger SQL注入LAB2 推荐一个测试人必备的Skills,从功能到性能全搞定(附详细实操和安装下载方式) 筑基期:掌握Odoo基础核心知识点02(Odoo XML 开发方式详解) GLM模型这么火,咱们用vllm也咧一个呗! 深入理解 AbortController:从底层原理到跨语言设计哲学 字符串学习笔记 多租户系统框架的基础模块设计和分析设计 Apache SeaTunnel Zeta 为什么能做到“又快又稳”? AI开发-python-LangGraph框架(3-26-LangGraph基本概念及第一个简单样例) Vue 3 组件通信,别只会用 Props 和 Emits 了,这几个狠活儿你得看看 ElasticSearch7.X版本配置密码 用Manim实现动态交点计算--从一个动点问题说起 团结引擎+Addressable+Instant Game打包抖音小游戏 function call 实战:让 LLM 自动判断 pod 异常、调用日志工具并完成故障分析 bubseek —— 让 Agent 的足迹,变成团队的洞察 通过 C# 读取并导出 PDF 书签 如何用 GitHub Actions 实现 Steam 自动化发布 【从0到1构建一个ClaudeAgent】并发-后台任务 .NET 高级开发 | 定制 ASP.NET Core 框架 电子小白:什么是运算放大器(运放) zero2Agent:面向大厂面试的 Agent 工程教程,从概念到生产的完整学习路线 堆上的ORW HC32F460 USB CDC通信异常:非对齐访问异常排查 20260413-Hyperbridge 攻击事件:发生在默克尔山上的验证绕过 那些喊着AI 要淘汰你的人,正在靠你的焦虑赚大钱! 深度学习进阶(八)Swin Transformer 最小二乘问题详解19:带先验约束的增量式SFM优化与实现 SnapTranslate 3.0 正式发布:全局划词翻译 + 完整英语学习闭环,一站式搞定查词、记词、复习 工作的意义、工作的困难认知再思考 .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 上下文工程是什么?过时了么?一文讲明白! - 一枫说码 开了 TUN 模式还是直连?90% 的人都踩过这个坑 AScript扩展多种脚本语言 - rockey627 AI 学习笔记:Agent 的记忆机制 你能被装进一个文件里吗?——7 万人把同事"蒸馏"成了 AI - 我没有三颗心脏 Claude Code 通关手册(七):给 AI 装上技能包——Skills 完全指南 - 暮色之狐 在浏览器中快速编辑代码:VSCode Web 集成实践 - Newbe36524 蒸馏自己 skill?基于 Deepseek 的蒸馏器,丐版蒸馏方式,简单便捷 - To_Carpe_Diem Spring AI Aliababa和AgentScope,哪个更好? - 苏三说技术
从零学习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 集群的稳定性。