尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Kafka核心设计与生产实践:从消息队列原理到避坑指南

Kafka核心设计与生产实践:从消息队列原理到避坑指南 先说下背景。我自己在业务系统里折腾消息中间件也有些年头了从最早用 RabbitMQ 做异步任务到后来在数据链路里大规模用 Kafka 承载日志和用户行为数据中间踩过的坑不少。这篇文章我不会按教科书的方式给你罗列概念而是从一个使用者的角度把 Kafka 最核心的设计思路、最容易踩的雷区以及怎么在实际项目里把它用好讲清楚。无论你是刚接触消息队列的新人还是已经用过却没深究过原理的开发者这篇文章都能帮你建立一套完整的认知框架。Kafka 本质上是一个分布式的消息队列但它的设计思路和传统消息队列比如 RabbitMQ、RocketMQ有很明显的区别。最直观的感受就是它更像一个可重放的高性能日志系统而不是一个即取即走的消息管道。理解了这句话后面所有概念都会变得顺理成章。1. 先搞清楚消息队列解决了什么问题Kafka 又特殊在哪1.1 从两个最常见的业务痛点说起第一类痛点是同步调用太慢。假设你有个订单系统用户下单后需要发送短信通知、赠送积分、更新推荐引擎的数据、记录行为日志。如果这些操作全部在订单请求链路里同步执行用户点击提交订单后可能要等上两三秒才能看到成功页面。引入消息队列后订单服务只需要把订单创建成功这个事件写入队列然后立刻返回后续的短信服务、积分服务各自订阅这个事件异步执行各自的任务接口响应时间直接从秒级降到毫秒级。第二类痛点是流量突增打垮数据库。典型的场景是秒杀活动一瞬间涌入的请求量可能是平时的一百倍。如果所有请求都直接打数据库数据库大概率会挂掉。用消息队列做削峰填谷后请求先进入队列后台消费者按自己的处理能力匀速处理。队列本质上是把瞬时高峰转换成了持续的压力给下游系统争取了缓冲时间。除了这两点消息队列还天然解决了系统解耦问题——生产数据的服务不需要知道谁在消费这些数据新增一个消费方也完全不需要改动已有代码。这就是为什么几乎每个中大型互联网项目里都能看到消息队列的身影。1.2 点对点模型和发布订阅模型Kafka 属于哪种理解消息队列的两种基础模型很重要因为很多概念上的混淆都源于此。点对点模型Point-to-Point里一条消息只能被一个消费者消费。消息进入队列后多个消费者争抢谁抢到谁消费消息一旦被消费就从队列里移除。这种模型适合任务分发比如一个任务队列被十个 worker 消费。发布订阅模型Pub/Sub里消息会广播给所有订阅了该主题的消费者。每个订阅者都能收到同一条消息的完整副本彼此之间互不影响。这种模型适合事件广播比如一条用户注册成功的事件要被短信服务、CRM 服务、数据统计服务同时处理。Kafka 本质上是一个发布订阅模型但它做了一层非常关键的扩展消息不会被立即删除而是按策略保留一段时间比如 7 天消费者可以按偏移量重放历史消息。这一点让它和传统 MQ 有了本质区别也是它敢说自己更像个日志系统的根本原因。1.3 为什么把 Kafka、RabbitMQ、RocketMQ 放在一起对比很多人在做技术选型时都会纠结这三者。我直接给结论再解释理由对比维度KafkaRabbitMQRocketMQ定位分布式消息流平台传统消息代理分布式消息中间件吞吐量极高百万级/秒中万级/秒高十万级/秒消息模型发布订阅为主支持消费者组点对点和发布订阅都支持发布订阅为主支持事务消息追溯按偏移量重放可读历史消费即删除支持按时间回溯可靠性通过副本机制保证支持多种确认机制极高支持事务消息学习成本概念多但严谨概念简单易上手概念适中RabbitMQ 的优势是功能全面、开箱即用路由规则灵活比如直连、主题、扇形交换器适合业务逻辑复杂、消息量不大的企业内部系统。RocketMQ 在阿里生态里用得非常多金融级的事务消息是一大亮点。Kafka 的强项则是极致的吞吐量和消息重放能力所以在数据管道、日志聚合、用户行为跟踪这些海量数据流动场景里几乎是公认的标准选择。2. Kafka 的核心设计为什么它能扛住每秒百万级消息2.1 追加写入与顺序读写的魔法数据库的 B 树索引在写入时需要进行随机 IO磁盘寻址是最大的性能瓶颈。Kafka 反其道而行之所有消息都只追加Append到分区日志文件的末尾写入操作永远是顺序的。顺序写磁盘的速度可以做到接近内存速度因为操作系统已经为此做了大量优化比如预读和延迟写入。消息持久化到磁盘后消费时同样利用顺序读的能力。消费者读取某个分区时按偏移量依次向后读取磁盘预读机制会把相邻的数据块一并加载到页缓存里下一次读取大概率直接命中缓存。这里我要提醒一点很多人误以为 Kafka 的高吞吐是因为它把数据放在内存里、不落盘。这是完全错误的认知。Kafka 恰恰是全部落盘却依然快快就快在这套顺序读写加页缓存的设计上。它利用了操作系统层面最基本的机制却取得了远超传统随机读写数据库的效果。2.2 分区Partition并行度的根本来源Kafka 的高吞吐另一个关键来自分区。一个主题Topic可以包含多个分区每个分区是一个有序的日志文件。生产者写消息时按规则比如按 key 哈希选择分区消费者读消息时每个分区对应一个消费者实例来读。分区的意义在于横向扩展对生产者来说多个分区意味着可以并行往多个文件里写总吞吐量可以线性扩展对消费者来说分区是消费并行度的上限——如果一个主题有 10 个分区那么最多可以有 10 个消费者实例同时消费每个实例负责一个分区。想加消费并发先加分区数。但是分区不是越多越好。每个分区都会对应一组文件和线程资源分区过多时文件句柄数、内存占用都会上升出故障时的恢复时间也会变长。实践经验是分区数建议控制在 broker 数量的 2~3 倍以内配合数据量综合评估。2.3 零拷贝与页缓存Kafka 性能的隐形翅膀传统的数据传输流程是磁盘 → 内核缓冲区 → 用户缓冲区 → Socket 缓冲区 → 网卡中间涉及多次拷贝和用户态/内核态切换。Kafka 利用了操作系统的零拷贝特性直接把磁盘文件数据从内核态通过 DMA 传输到网卡中间跳过了用户态。发送端的实际效果是数据从生产到消费全程基本不经过应用进程的内存吞吐自然高。再加上页缓存这一层Kafka 不像传统消息队列那样维护一套复杂的消息在内存还是磁盘的状态机。所有数据统一走磁盘读的时候由操作系统页缓存兜底省掉了大量 JVM 内存管理开销。这也是为什么 Kafka 服务端可以做到堆内存占用很稳定不会像某些中间件那样动不动就 OOM。2.4 消息保留策略不删除才是设计精髓传统 MQ 消费完消息就删除Kafka 却默认保留 7 天可配置甚至可配置永久保留。这个设计初看有点反直觉仔细想非常巧妙支持离线/故障恢复。消费者挂了 3 天恢复后可以从上次提交的偏移量接着读一条消息都不丢支持业务回溯分析。数据团队可以从头重新消费一遍历史数据用于修正之前的计算错误支持流处理。Kafka Streams 等流计算框架正是利用这种可重放的能力来实现精确一次的语义。2.5 一个类比Kafka 像是可以倒带的磁带说了这么多理论用一个生活化的类比收一下Kafka 的每个分区类似于一盘磁带写入时磁头永远往后移动读取时可以随时把磁头拨回卷头重放。磁带有多少音轨分区决定了能同时有多少人操作这盘磁带磁带本身录了什么不会因为被听过就消失保留多久由你自己的策略决定。这套磁带哲学贯彻了 Kafka 的所有核心机制。3. 必须吃透的几个基础概念分区、副本、消费者组和偏移量3.1 Topic、Partition、Broker 与 Record 的对应关系我建议在学习任何消息队列时都先建立一套名词与实物的对应关系否则后面读文档会非常痛苦。Broker一台 Kafka 服务器就是一个 Broker。一个集群由多个 Broker 组成每个 Broker 有自己的编号broker.id负责存储部分分区的数据。Topic逻辑上的消息分类。每个 Topic 可以理解为一个目录目录里有多份数据文件。给业务起 Topic 名时建议用系统.事件.类型的格式比如order.created.v1方便管理。PartitionTopic 下面的物理分片。一个 Topic 的数据根据分区策略分布到多个分区中每个分区内部有序。Record一条实际的消息。包含 key、value、timestamp 和 headers。key 用于分区路由value 是真正的业务负载timestamp 标记消息产生时间可以影响日志保留策略。Offset消息在分区内的顺序编号。从 0 开始严格递增描述了这条消息在这个分区里排第几。一定要记住Kafka 保证顺序性的粒度是分区而不是整个 Topic。你只能保证同一个分区内的消息有序跨分区有序是做不到的也不需要做到。3.2 副本机制与 ISRKafka 高可用的底气每个分区都可以配置多个副本副本之间是一主多从的关系。主副本Leader负责处理读写请求从副本Follower负责同步数据主副本挂了之后从副本能顶上。这里有一个重要的概念叫 ISRIn-Sync Replica同步副本集合。它表示当前与主副本保持同步的副本列表。Kafka 不会把领导权交给一个数据严重落后的从副本只会从 ISR 里选举新的 Leader这样能最大程度避免数据丢失。为什么 ISR 概念容易被面试官追问因为理解了 ISR就理解了 Kafka 在性能和一致性之间的取舍。如果要求所有副本都确认写入才算成功延迟会很高但数据最安全如果只要 Leader 确认就返回性能最好但极端情况下可能丢数据。Kafka 可以通过acks参数在这个光谱上自由调整下面第 5 节会细说。我在生产环境里遇到过某个副本频繁同步超时被踢出 ISR 的情况典型原因是 Broker CPU 打满或者网络抖动。这种时候业务可能没有任何感知但集群的冗余能力其实已经悄悄下降了需要监控 ISR 收缩情况才能及时发现。3.3 消费者组一条消息怎么被多个业务系统消费消费者组Consumer Group是 Kafka 实现发布订阅的关键机制。当一个 Consumer 实例加入一个消费者组时组内所有实例会共同分担该组订阅的主题里的全部分区。Kafka 内部的协调器Group Coordinator负责分配哪个消费者处理哪些分区。分配完后每个分区只被组内一个消费者消费——这是 Kafka 实现消息只处理一次但不丢失的基础。多个消费者组订阅同一个 Topic 时每个组都能收到全量消息。比如你有个order.created的主题短信组读它发短信积分组读它加积分两组互不干扰。一个关键公式再强调一遍消费并发上限 分区数。如果只有一个消费者但主题有 10 个分区它要处理全部 10 个分区的数据如果消费者加到 11 个第 11 个只会白白闲置不会比 10 个消费者更快。3.4 偏移量提交自动提交还是手动提交偏移量Offset记录的是消费者当前读到了哪个位置。消费者每消费一条消息偏移量就可以往前推进。偏移量有两种提交方式这也是重复消费和消息丢失的分水岭自动提交enable.auto.committrue每隔auto.commit.interval.ms默认 5 秒自动提交当前消费到的偏移量。好处是省事坏处是消费逻辑执行了一半/刚拉取到一批消息但还没处理完时偏移量已经提交了——如果消费者这时候宕机重启就相当于这批消息被跳过了。这就是消息丢失的一大来源。手动提交enable.auto.commitfalse业务代码在处理完消息后主动调用 commitSync() 或 commitAsync()。好处是确认处理成功才提交偏移量保证不丢消息。坏处是要自己管好幂等性和异常处理。我的推荐很明确生产环境一律使用手动提交理由在第 5 节会结合具体案例展开。4. 动手环节从零搭一个 Kafka 环境并跑通生产消费4.1 版本选择为什么建议直接上 KRaft 模式的 3.x早期的 Kafka 依赖 ZooKeeper 做集群元数据管理Broker 注册、Topic 元信息、Controller 选举都靠它部署时除了要搭 Kafka 集群还要额外搭一套 ZK 集群维护成本不低而且 ZK 本身也容易成为瓶颈。从 Kafka 3.0 开始官方逐步推进 KRaft 模式内置的元数据管理机制去除 ZooKeeper 依赖到 3.3 版本之后 KRaft 已经可以在生产环境使用了。我的建议是如果你是新手或小规模集群直接选择 3.5 以上的 KRaft 版本只有一个组件配置简单适合快速验证和学习如果你所在团队有成熟的 ZK 运维体系旧版本升级时再慢慢过渡不要强行迁移。4.2 Windows 环境下安装并启动 Kafka含验证很多人习惯在 Windows 上做学习验证。虽然 Kafka 官方没有提供 Windows 原生启动脚本但通过内置的 Windows 批处理脚本跑起来并不难。下面是我在 Windows 11 上实测无误的步骤第一步安装 JDK建议 OpenJDK 11 或 17配置好 JAVA_HOME 环境变量。第二步去 Apache Kafka 官网下载二进制压缩包比如 kafka_2.13-3.6.2.tgz用任意解压工具解压到类似D:\kafka的目录。第三步因为 KRaft 模式需要先生成集群 ID 并格式化存储目录打开命令行进入 Kafka 目录执行bin\windows\kafka-storage.bat random-uuid把输出的 UUID 记录下来然后用它格式化存储目录bin\windows\kafka-storage.bat format -t 刚拿到的UUID -c config\kraft\server.properties第四步启动 Kafka 服务bin\windows\kafka-server-start.bat config\kraft\server.properties看到日志输出started (kafka.server.KafkaRaftServer)就表示启动成功了。KRaft 模式默认监听 9092 端口同时会有一个内部控制器端口配置注意不要改错配置里的 listeners。第五步验证。再开两个命令行窗口一个窗口创建并启动生产者bin\windows\kafka-console-producer.bat --bootstrap-server localhost:9092 --topic test-topic另一个窗口启动消费者bin\windows\kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic test-topic --from-beginning在生产者窗口输入任意文本回车后消费者窗口应该能立刻收到同样的内容。到这一步你的本地 Kafka 就算正式跑起来了。这里分享一个 Windows 上容易踩的坑路径里尽量不要有中文和空格如果监听 localhost 连不上检查一下server.properties里的listeners是不是配成了PLAINTEXT://localhost:9092而不是只配了内网 IP。4.3 用 Java 代码跑通生产和消费跑通命令行之后再用代码写一个生产者和消费者你会对 Kafka 的使用模式有更直观的感觉。这里用 Maven 项目引入官方客户端依赖dependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version3.6.2/version /dependency生产者代码Properties props new Properties(); props.put(bootstrap.servers, localhost:9092); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(acks, all); KafkaProducerString, String producer new KafkaProducer(props); for (int i 0; i 100; i) { String key user- (i % 10); String value {\userId\: i ,\event\:\click\}; producer.send(new ProducerRecord(order.created.v1, key, value)); } producer.close();消费者代码Properties props new Properties(); props.put(bootstrap.servers, localhost:9092); props.put(group.id, demo-group); props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(enable.auto.commit, false); KafkaConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(Collections.singletonList(order.created.v1)); while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(100)); for (ConsumerRecordString, String record : records) { System.out.println(收到: partition record.partition() , offset record.offset() , key record.key() , value record.value()); } consumer.commitSync(); }看到这里你应该对生产者在分区写入、消费者在分区读取并提交偏移量有了具象的感知。如果还有余力试试给消息加一个key你会发现相同 key 的消息会稳定落到同一个分区——这是后续解决顺序性问题的基础能力。4.4 聊聊纯软件学习者可能会好奇的 Windows 消息队列搜索热词里有 msmq 和 Windows 消息队列这其实和 Kafka 是两类东西。MSMQ 是微软提供的操作系统级消息队列组件集成在 Windows 系统里常用于较早的 .NET 应用做分布式通信。它胜在与 Windows 生态无缝集成、部署简单但功能远没有 Kafka 丰富吞吐量和跨平台能力也有明显天花板。如果你现在做的新项目没有历史包袱我个人不建议选 MSMQ如果是在维护老系统理解它主要是为了兼容和改造迁移。5. 生产环境避坑不丢消息与重复消费问题怎么解5.1 三个不丢消息的环节一个都不能少消息不丢是一个系统性要求不是靠某一个参数就能保证的。按生产链路拆开来看至少要保证三个环节生产端不丢。生产者发送消息时如果设置acks0意味着发出去就不管结果完全没有确认吞吐最高但也最容易丢。设置acksall加上min.insync.replicas2表示消息要写入主副本且至少一个同步副本确认后才返回成功。理论上只要不是所有同步副本同时宕机消息就不会丢。这是我最推荐的生产配置。Broker 端不丢。消息写入 Broker 的分区日志文件后Kafka 会通过副本机制同步。但如果只有一个副本Broker 磁盘坏了数据就真的没了。答案是用副本因子replication.factor给重要 Topic 配 3 个副本同时设置unclean.leader.election.enablefalse避免那些数据落后的副本被选为主副本这会直接导致历史数据丢失。消费端不丢。消费者取到一批消息后如果还没处理完就崩溃偏移量还停留在旧位置重启后会重新拉取这些消息。看起来重消费好像不是丢但如果你启用了自动提交情况会反过来一批消息刚拉取到本地还没来得及处理偏移量因为定时器触发已经提交了下一条的位置重新启动后这批消息就永远不消费了。这正是消费端丢消息最常见的坑解法就是手动提交且务必在业务处理成功之后再提交。5.2 重复消费不可避免怎么用幂等性收尾即使你做对了上面所有配置分布式系统里依然会重复消费消费者处理完消息后在提交偏移量前崩溃了重启后重新拉取到同一条消息又处理了一遍。这不是配置问题而是至少一次at-least-once语义里的固有特征——只要你想保证不丢就必须接受潜在的重复。应对重复的核心思路是幂等消费无论同一条消息被处理多少次结果都一样。实际操作上有几种模式唯一键去重。业务表里加一个唯一索引消息里带唯一事件 ID比如订单号事件类型时间戳插入时用 insert ignore冲突就不处理。这是最朴素也最常用的方式。状态机校验。如果业务流程有明确的状态流转消费时先检查当前状态是否允许执行该动作不允许就直接跳过。比如已经支付的订单收到关闭订单消息时可以直接忽略。数据库乐观锁。更新数据时带上版本号更新行数影响为 0 就说明本次操作已过期。我见过不少团队把精力花在怎么才能不重复发消息上其实方向偏了。重复消费是分布式系统里的常态重点是让下游扛住重复而不是徒劳地消灭重复。5.3 顺序性需求的三种常见场景与解法Kafka 只能保证分区内有序。如果你的业务对消息顺序有严格要求按场景分情况解决场景一同一个业务对象的操作必须有序比如订单先创建、后支付、再关闭。解法用业务对象 ID 作为消息 key相同 key 会路由到同一个分区分区内天然有序。这是最简单有效的方式。场景二多个主题之间存在先后依赖比如用户信息变更事件和用户行为分析事件前者必须先到。解法不要让这类有依赖关系的事件分散到不同主题尽量合并到同一主题同一 key如果无法合并考虑在下游消费端按时间戳做等待/排序处理成本较高尽量避免设计成这种结构。场景三需要对消息进行全局严格排序。Kafka 本质上实现不了跨分区全局有序除非你把主题分区数设为 1。但分区分成 1 意味着并发度归零吞吐大打折扣。我建议先反问业务这个全局有序是真实需求还是可以通过业务层对账/重排序解决绝大多数情况都不需要真正的全局有序。5.4 面试题角度消息队列重复消费与积压问题热词里的kafka面试题及答案说明很多人是为了面试准备的那就总结一下最常见的追问链路面试官通常会问消息队列有什么作用你答削峰、解耦、异步。深入一点会问Kafka 怎么保证消息不丢你沿生产者 acks、Broker 副本、消费者手动提交三个维度回答。再深入一点怎么保证不重复消费你答至少一次语义下无法避免但可以靠幂等设计兜底。如果面试官追问那 Kafka 能不能保证消息顺序你就可以把上面三种场景和 key 路由机制讲一遍。这套链路里最容易被翻车的地方就是只答对一半概念却没有自己的实战案例。任何一层如果能补一个实际项目里的经验比如我之前因自动提交丢过数据改成手动提交后消费者吞吐下降了但可靠性完全不一样了效果会立刻不同。6. 日常观测与性能调优入口可视化工具和延迟排查6.1 别再问Kafka 有没有 UI 界面了这些工具够你选Kafka 的庞大配置和分区状态全在配置文件里如果没有可视化工具排查问题时全靠命令行和日志非常不直观。目前社区里常用、我实际用下来也靠谱的可视化工具有这几个工具定位优点注意点Kafka UIkafka-ui一站式 Web 管理支持多集群、查看 Topic/分区/消费者组、消息检索、动态调整配置Java 服务部署一个实例即可Kafdrop轻量 Web 工具简单直接适合快速查看 Topic 和消息内容功能偏少不支持消费者组管理Offset Explorer原 Kafka Tool桌面客户端看分区、偏移量、消费者 lag 非常直观只适合本机使用不适合团队共享云厂商控制台托管版专属监控告警非常完备自建集群用不了我自己最常用的组合是 Kafka UI 做日常管理、Offset Explorer 查消费积压明细。注意一点这些工具本身都不参与消息流转只读元数据和消息内容可以放心部署。6.2 消息延迟高从哪个环节慢开始查搜索热词里的kafka消息延迟高是人人都想解决的问题。消息延迟的链路可以简化成生产者发送耗时 → 网络传输 → Broker 落盘耗时 → 消费者拉取间隔 → 消费者处理耗时 → 下次心跳/拉取的等待。我的排查顺序是第一步先看消费端有没有积压。用工具查看消费者组的 lag最后一条消息的偏移量和已完成偏移量的差值如果 lag 持续增大说明消费能力跟不上生产速度问题通常在消费端。第二步看消费端单条消息处理耗时。如果在消息处理逻辑里调用了慢数据库或慢接口poll 拉取下一个批次的周期会被拉长。此时用max.poll.records控制单批次条数、控制max.poll.interval.ms防止消费者因处理太久被判定为死亡以及考虑给消费组扩容前提是有足够分区。第三步看生产端有没有破洞。如果某个 Topic 的瞬时峰值极大而分区数不够生产端的吞吐就会卡在分区上限上。此时可以先确认batch.size和linger.ms的配合情况小幅调优缺分区就增加分区数和 Broker 数量。第四步看 Broker 本身。磁盘 IO 使用率、网络带宽、页缓存命中率都要纳入检查。Kafka 对磁盘吞吐要求很高机械硬盘或低性能云盘在高负载下最容易成为瓶颈。有条件直接用 NVMe 固态盘延迟表现会非常稳定。6.3 单条消息 1MB默认配置的隐藏坑搜索热词里的kafka 接收1m我理解是在说 1MB 左右的大消息能不能传给 Kafka。默认情况下Kafka 的单条消息最大是 1MBmessage.max.bytes这个配置项还关联了 broker 端、topic 端、生产者端、消费者端至少四个位置的参数。如果你只改了 broker 端而客户端没改生产或消费时会出现消息过大被拒绝的报错。正确修改方式是同时改动四个位置Broker 的message.max.bytes、Topic 的max.message.bytes、生产者的max.request.size、消费者的fetch.max.partition.bytes。比如要支持 5MB 消息# broker 端 server.properties message.max.bytes5242880 replica.fetch.max.bytes5242880// 生产者端 props.put(max.request.size, 5242880); // 消费者端 props.put(fetch.max.partition.bytes, 5242880);但我必须提醒大消息尽量不要走 Kafka。Kafka 的强项是海量小消息的高吞吐。单条 5MB 的消息会占用大量内存和网络资源还会延长落盘和复制时间压缩效果也很差。如果确实有大规模文件传输需求更合理的做法是把文件放到对象存储里Kafka 只传文件地址和元信息消费方再自行下载。6.4 从入门到能干活你的下一步这篇文章能帮你把 Kafka 的骨架搭起来但真正的掌握还是要靠动手。我的建议是接着做三件事第一用 Kafka UI 观察你本地创建的主题。创建一个有 3 个分区、副本因子 1 的主题启动两个消费者跑同一个组再看每个消费者各领到了哪些分区——你会立刻理解消费者组的分配模型。第二故意制造一次消费失败。在消费端处理消息时手动抛异常不用手动提交偏移量观察这个消息是不是会反复出现在日志里。这个印象比你读十篇文档都深刻。第三用 docker-compose 搭一个 3 节点 KRaft 模式的 Kafka 集群调整副本因子然后 kill 掉一个节点观察集群自动切换 Leader 的过程。高可用的概念只有亲手测过才踏实。最后分享一个个人经验。我刚开始学习 Kafka 时觉得它概念又多又绕ZooKeeper、ISR、消费者组每个名词都像一座山。后来我换了思路不去背概念而是在一台机器上反复生产、消费、kill、改配置、看监控一边操作一边对照文档里那些抽象名词。大概两周时间那些名词就自然而然长在脑子里了。Kafka 这种基础设施软件光看文档是不可能学会的把它跑起来亲手弄坏再修好——这才是最有效的入门路径。
返回列表