ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Kafka高可用集群部署与调优实战:从零搭建到故障演练

Kafka高可用集群部署与调优实战:从零搭建到故障演练 做 Java 后端的人和 Kafka 打交道基本是躲不开的。但说句实话单机跑通 Kafka 和把它作为高可用消息队列集群部署到生产环境完全是两码事。我见过不少项目开发环境用 Docker 起一个 Kafka 就完事等压测或者线上流量一上来Leader 切换、消息积压、消费延迟各种问题全冒出来。这篇文章把我自己从零部署三节点 Kafka 集群、做故障演练、再逐步调优的完整过程记录下来包含每一步的配置依据和踩过的坑适合准备在生产环境搭建高可用 Kafka 的 Java 后端同学参考。1. 为什么 Kafka 高可用不能靠多装几个节点了事1.1 高可用不等于副本数多先理解 Kafka 的可用性边界很多人以为高可用就是把副本因子设置成 3多复制几份数据就万事大吉。实际上 Kafka 的高可用依赖的是一个完整链路Broker 节点要能正常通信、Controller 要能快速完成 Leader 迁移、ISR 列表里的副本要跟上写入、客户端还要配置合理的 acks 和重试策略。任何一个环节掉链子集群都可能出现消息丢失或长时间不可用。举个例子我遇到过一个小型生产集群副本因子确实配置成了 3但 Producer 端一直在用默认的 acks0也就是发出去就不管了。某次一台 Broker 突然宕机数据其实已经写入 Leader 分区但因为是异步发送客户端本地并不知道成功还是失败等 Leader 切换完成后一批消息就彻底没了。事后排查发现服务端高可用配置得再完善客户端不配合照样丢数据。所以正确理解是高可用是服务端和客户端共同协作的结果。服务端负责通过分区副本和 ISR 机制保证数据的冗余与可恢复性客户端负责通过 acks、重试、幂等等参数确保数据在语义上真正写入成功。1.2 从 CAP 视角看 Kafka分区副本与 ISR 的取舍逻辑Kafka 本质上是一个基于日志的分布式系统它在 CAP 里并不是简单的AP或CP而是可以借助配置在两者之间滑动。默认情况下Kafka 允许每个分区在 Leader 不可用时从 ISRIn-Sync Replicas同步副本中选举新的 LeaderISR 里的副本都在持续拉取 Leader 的数据。只要 ISR 里有副本活着系统就能继续提供读写服务这时候更偏向可用性。但只要把unclean.leader.election.enable设置为 false并让min.insync.replicas大于 1Kafka 就会拒绝在 ISR 没有足够副本时选举落后太多的副本作为 Leader。这样做的代价是极端情况下比如 3 副本只剩 1 个 ISR而你又要求至少 2 个 ISR 才能继续写写入会被拒绝系统暂时不可用但数据不会丢。这就引出一个非常核心的决策你的业务到底能不能接受短暂的写入不可用如果能接受那就用min.insync.replicas2acksallunclean.leader.election.enablefalse这是数据零丢失的经典组合。如果不能接受短暂的不可用那就要放宽副本同步要求但必须接受极端情况下丢数据的风险。很多团队在这个选择上摇摆不定最后配置改成了四不像既不能保证可靠也不能保证可用。1.3 单机伪集群与生产集群的差距开发环境常见的做法是在一台机器上起三个 Broker端口不同、日志目录不同看起来是个集群但生产环境和它的差距非常大。单机伪集群的所有 Broker 共享同一台机器的 CPU、内存、磁盘 IO 和带宽宕机时实际上整台机器都不行了副本再多也只是数据在同一个物理机上复制没有任何故障域隔离。生产集群至少要满足三个条件Broker 分布在不同的物理机或至少不同的故障域上每个分区的副本尽量分散到不同机器客户端连接能够感知多节点而不是写死单点。很多人部署的时候只是把 server.properties 里的 broker.id 改一改然后复制三份以为就成了高可用集群。真的发生物理机宕机所有 Broker 一起挂高可用就成了笑话。2. 部署前必须拍板的决策点KRaft 还是 ZooKeeper、磁盘怎么选、JVM 给多大2.1 KRaft 与 ZooKeeper 模式元数据管理的迁移现状新部署 Kafka 集群第一个要决策的就是元数据管理模式。Kafka 早期依赖 ZooKeeper 保存 Broker、Topic、分区、Controller 等元数据同时也用 ZK 做 Leader 选举。ZooKeeper 本身也是一套需要高可用部署的系统等于变相增加了运维成本。我们团队早期维护的 Kafka 1.x 集群ZooKeeper 经常出现会话超时、节点数据不一致等问题排查起来非常痛苦。Kafka 2.8 引入了 KRaft 模式简单理解就是用 Kafka 自己的内部日志来管理元数据去掉 ZooKeeper 依赖。从 Kafka 3.3 开始 KRaft 可以用于生产。如果你是全新搭建集群我建议直接用 KRaft 模式少维护一套 ZooKeeperController 的高可用由 Kafka 自己保证。但要注意KRaft 模式下需要先格式化存储目录而且 Controller 节点的node.id、controller.quorum.voters等配置要和 Broker 保持一致否则集群起不来。混合模式Controller 和 Broker 在同一个进程里适合中小型集群独立 Controller 模式适合大规模集群。我们的三节点集群用的是混合模式每个节点既是 Broker 也是 Controller既减少了机器数量又能保证元数据高可用。2.2 磁盘选型与文件系统挂载页缓存之外的硬件变量Kafka 官方文档一直在强调顺序读写的优势说机械硬盘也能跑出不错的效果。这句话本身没错但前提是顺序。生产环境一旦出现分区重平衡、多个 Topic 同时写入、Consumer 回溯消费磁盘上的 IO 模式立刻变得杂乱机械硬盘的性能会大幅跳水。我们第一次压测时用的就是普通 SATA 机械盘单 Topic 单分区写吞吐还能看跑 12 个分区后直接跌到不足原来的三分之一。磁盘挂载也有讲究。生产环境建议用 xfs 文件系统挂载参数加上noatime避免访问时间戳的频繁写入。Kafka 的日志目录通过log.dirs配置你可以配置多个目录比如/data1/kafka,/data2/kafkaKafka 会把不同分区的日志目录分散到这些目录上减轻单盘压力。不要把这些目录放在系统盘上系统日志和 Kafka 日志抢 IO 的后果很严重。2.3 JVM 堆与操作系统参数Kafka 是 Java 应用但不全靠堆Kafka 本身是一个 Java 应用但它对堆内存的依赖远没有你想象的大。Kafka 高性能的核心之一是页缓存Page Cache读写消息时优先走操作系统缓存。这也意味着你给机器配置的内存里要预留足够大的空间给页缓存而不是一股脑全分给 JVM 堆。我见过有同事把 Kafka 的KAFKA_HEAP_OPTS直接设成-Xmx16G结果频繁发生 Full GC吞吐反而不如之前 4G 堆的时候。原因就是堆太大GC 扫描时间长而且留给页缓存的内存变少了。对于绝大多数 Kafka 节点堆内存给 4GB 到 6GB 就足够了重点优化的是 GC 策略。我们用的配置是-Xms4g -Xmx4g -XX:UseG1GC -XX:MaxGCPauseMillis20实测停顿控制得不错。操作系统层面有两个参数我每次部署都会检查文件描述符上限和vm.swappiness。Kafka 会打开大量文件句柄ulimit -n建议调到 100000 以上。vm.swappiness默认是 60在 Kafka 机器上建议调低到 10 左右避免不必要的交换分区写入。2.4 集群节点规划奇数节点、机架感知与分区分配节点数量上最常见的生产配置是 3 节点或 5 节点。如果使用 KRaft 模式Controller 节点数量建议是奇数方便 quorum 选举。3 节点可以容忍 1 个节点宕机5 节点可以容忍 2 个节点宕机。对于大多数中大型业务3 到 5 个 Broker 节点已经足够。副本因子和分区数的规划容易走极端。副本因子建议生产环境至少 3这样配合min.insync.replicas2才能做到写入层面的高可用。分区数不是越多越好分区太多会导致每个分区数据量过小、文件碎片变多还会增加 Controller 和客户端的管理开销。一个经验法则是先估算目标吞吐参考单分区实测吞吐分区数等于目标吞吐 / 单分区吞吐再往上冗余 20% 到 30%。举个具体例子如果单分区能扛住 20MB/s目标吞吐是 200MB/s初始分区数可以定在 12 到 15。如果机器分布在多个机架或者多个可用区尽量开启机架感知rack awareness把副本分配在不同机架上这样即使一个机架断电其他副本依然能提供服务。3. 从零部署一个三节点 Kafka 集群操作步骤与隐藏细节3.1 节点配置与二进制包准备这里以 Kafka 3.6 版本为例三台 Linux 服务器每台配置 8C16G操作系统 CentOS 7.9JDK 使用 OpenJDK 17。Kafka 3.x 要求 JDK 8 以上但高版本 JDK 在性能和安全上更省心。下载二进制包时要注意到 Apache Kafka 官网下载kafka_2.13-3.6.2.tgz这种命名格式的包里面包含了 Broker 和自带工具不需要额外装 ZooKeeperKRaft 模式。下载后解压到/opt/kafka创建数据目录和日志目录tar -xzf kafka_2.13-3.6.2.tgz -C /opt mv /opt/kafka_2.13-3.6.2 /opt/kafka mkdir -p /data/kafka-logs mkdir -p /data/kafka-metakafka-logs存放消息日志kafka-meta存放 KRaft 元数据两个目录要分开方便后续排查和备份。3.2 核心配置项逐行说明每台节点的config/server.properties都要根据角色做修改。以 node-1 为例核心配置如下process.rolesbroker,controller node.id1 controller.quorum.voters1192.168.1.11:9093,2192.168.1.12:9093,3192.168.1.13:9093 listenersPLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 advertised.listenersPLAINTEXT://192.168.1.11:9092 controller.listener.namesCONTROLLER inter.broker.listener.namePLAINTEXT log.dirs/data/kafka-logs num.partitions3 default.replication.factor3 min.insync.replicas2 log.retention.hours72process.rolesbroker,controller意味着这个节点同时承担 Broker 和 Controller 角色。混合模式下三台节点都一样。controller.quorum.voters是 KRaft 模式最关键的一行它列出了所有 Controller 节点的 ID 和地址三台机器要填写完全相同的内容只是各自node.id不同。advertised.listeners是给客户端和集群其他节点看的地址这里最容易踩坑如果不填或者填了localhost客户端在另一台机器上就会因为拿到错地址连不上。3.3 启动、验证与 systemd 托管KRaft 模式首次启动前需要生成集群 ID 并格式化存储目录这个步骤只在第一次做而且必须三台节点使用同一个 Cluster ID# 在 node-1 上生成 Cluster ID /opt/kafka/bin/kafka-storage.sh random-uuid # 格式化每台节点执行用同一个 Cluster ID /opt/kafka/bin/kafka-storage.sh format -t Cluster-ID -c /opt/kafka/config/server.properties格式化会清空元数据目录所以生产环境一定要确认目录里没有重要数据再执行。格式化完成后就可以启动了/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties为了生产稳定性建议用 systemd 托管 Kafka 进程。创建一个/etc/systemd/system/kafka.service文件ExecStart 指向启动脚本然后systemctl enable --now kafka。托管的好处是机器意外重启后 Kafka 能自动拉起不用人为干预。验证集群是否正常最直接的方式是用自带的工具查看 Broker 列表/opt/kafka/bin/kafka-metadata-quorum.sh --bootstrap-server 192.168.1.11:9092,192.168.1.12:9092,192.168.1.13:9092 describe --status能看到三个节点都处于Leader或Follower状态说明 KRaft 元数据集群已经就绪。3.4 创建带多副本的 Topic 并验证高可用集群跑起来后创建一个测试 Topic指定 3 个副本/opt/kafka/bin/kafka-topics.sh --bootstrap-server 192.168.1.11:9092,192.168.1.12:9092,192.168.1.13:9092 \ --create --topic test-ha --partitions 3 --replication-factor 3然后查看分区副本分布/opt/kafka/bin/kafka-topics.sh --bootstrap-server 192.168.1.11:9092,192.168.1.12:9092,192.168.1.13:9092 \ --describe --topic test-ha正常情况下每个分区的Replicas会分布在 3 台节点上Isr也是 3 个。如果某个分区的 ISR 数量长期少于副本数说明有节点同步落后或者已经宕机要尽快排查。我习惯在创建完 Topic 后立刻做一次简单的生产消费测试用kafka-console-producer.sh写入几条消息再用kafka-console-consumer.sh从--from-beginning消费出来。这条链路通了集群的基本读写能力才算验证通过。4. 故障演练实录杀掉 Leader、重启 Broker集群到底发生了什么4.1 模拟 Broker 宕机消费者的感知时间集群部署不是跑起来就算了高可用能力必须通过故障演练验证。我第一次做演练时直接在 node-1 上执行kill -9模拟最极端的宕机场景然后观察消费者端的情况。当时消费者是一个 Java 应用使用默认的session.timeout.ms45000新版本默认值旧版本是 10000。节点 kill 掉之后消费者组并不会立刻感知而是等到 Session 超时后才会触发 Rebalance重新分配分区。也就是说Leader 已经切换了但消费者的感知延迟可能长达几十秒。对很多业务来说这个时间不可接受。通过这次演练我认识到两个问题一是 Broker 端要配置controller.quorum.election.timeout.ms适当调小让 Controller 快速完成新 Leader 选举二是消费端要根据业务容忍度调小session.timeout.ms和heartbeat.interval.ms比如分别设为 10000 和 3000让消费者更快感知 Broker 变化并触发 Rebalance。但也不能把session.timeout.ms调得太小否则网络抖动会导致消费者频繁掉线。4.2 数据会不会丢acks、min.insync.replicas 与 unclean.leader.election演练宕机的同时我还在持续发送消息重点观察有没有丢数据。这里要分清两个层面Broker 集群层面的数据冗余和客户端写入语义层面的可靠性。Broker 层面Leader 节点宕机后Controller 会从 ISR 中选一个新 Leader。如果使用unclean.leader.election.enablefalse那么只有 ISR 里的副本才有资格成为 Leader没跟上数据的副本不会上位这样不会丢已经提交的消息。反过来如果这个参数设置成 trueOSROut-of-Sync Replica副本也可能被选为 Leader它缺少的数据就永久丢了。客户端层面如果 Producer 设置acksall并且 Topic 设置了min.insync.replicas2那么写入必须被至少 2 个节点确认才算成功。这意味着在 3 副本集群中即使 1 台节点完全宕机Producer 依然可以继续写入但如果同时挂掉 2 台节点ISR 里只剩 1 个副本min.insync.replicas2的要求无法满足写入会被拒绝从而保证不会把数据写进一个不可靠的单副本环境里。我用这个配置组合做了三次连续 kill 演练第一次 kill 一个节点生产消费都正常只有短暂的 Leader 切换第二次 kill 两个节点Producer 开始抛NotEnoughReplicasException这正是预期行为保证不丢消息等节点恢复后写入自动恢复。这个结果让我对高可用有了更直观的认识高可用不是不出问题而是问题发生时数据安全有底线。4.3 恢复过程中的踩坑磁盘写满与 Controller 迁移故障演练还暴露出一个运维坑宕机节点恢复后如果它上面的副本落后太多会触发大规模的副本追赶和分区重平衡。我曾经在一次演练后看到三台机器的磁盘占用率迅速飙升差点写满。原因是被宕机节点重新加入集群后所有需要追赶的副本会同时从 Leader 拉数据期间还会产生大量临时文件。解决办法有两个一是平时监控磁盘水位给日志目录预留至少 20% 的空闲空间二是在节点恢复后如果副本追赶压力大可以临时设置replica.fetch.max.bytes和副本限流参数控制副本拉取速度让 IO 压力平滑释放。还有就是 Controller 迁移。如果宕机的是当前 Controller 节点其他节点会重新选举 Controller。这个动作通常很快但如果在高负载下选举过程可能被延迟。我在一次演练中观察到一个现象Broker 已经恢复但 Topic 的 Leader 没有马上恢复查了一圈发现是元数据缓存还没有刷新。重启两个节点后消失。后来养成了一个习惯节点恢复后用kafka-leader-election.sh手动触发一次优选的 Leader 选举让 Leader 尽可能均匀分布回原有的节点上。5. 生产环境性能调优先压测再调参别被默认配置骗了5.1 用 Kafka 自带工具做基准压测默认配置能跑但离跑得好还有很大距离。调优前必须先做基准压测搞清楚当前集群的瓶颈在哪里。Kafka 自带两个压测脚本分别是kafka-producer-perf-test.sh和kafka-consumer-perf-test.sh。压测 Producer 时我常用的命令/opt/kafka/bin/kafka-producer-perf-test.sh \ --topic perf-test \ --num-records 1000000 \ --record-size 1024 \ --throughput -1 \ --producer-props bootstrap.servers192.168.1.11:9092,192.168.1.12:9092,192.168.1.13:9092 \ acksall \ buffer.memory67108864 \ batch.size32768 \ linger.ms20 \ compression.typelz4--throughput -1表示不限制吞吐让客户端全力压。脚本最后会输出吞吐量records/sec、MB/sec和延迟百分位。如果第一轮压测吞吐远低于预期不要急着调 Producer 参数先看 Kafka Broker 所在机器的 CPU、网络和磁盘 IO 是不是已经被打满。5.2 Producer 端核心参数acks、batch、linger.ms、压缩Producer 端参数对吞吐影响最直接的有四个acks、batch.size、linger.ms和compression.type。acksall是可靠性要求但它会降低吞吐因为需要等所有 ISR 确认。如果业务允许少量数据丢失用acks1吞吐会高不少但我不建议为了性能牺牲可靠性性能可以通过其他参数找回来。batch.size控制 Producer 每个批次最多攒多少字节再发送默认是 16KB。如果消息平均大小只有几百字节这个值可以适当调大比如 32KB 或 64KB减少网络请求次数。linger.ms是批次在内存里等待更多消息加入的时间默认是 0也就是有消息就发。调大linger.ms到 10ms 或 20ms 能显著提升批处理效率代价是增加少量人为延迟。如果业务要求低延迟linger.ms不宜超过 5ms。compression.type是很容易被忽略的配置。开启压缩不仅能降低网络带宽占用还能提高吞吐。在 CPU 有富余的机器上建议 Producer 端开启lz4或zstd压缩。我们实测中zstd压缩率最好但 CPU 占用稍高lz4则是在压缩率和 CPU 消耗之间比较均衡的选择。消息体是 JSON 文本的场景压测吞吐可以提升 30% 到 50%。5.3 Consumer 端核心参数fetch、max.poll.records 与消费 LagConsumer 端调优的目标通常不是吞吐而是消费延迟和稳定性。最容易出问题的参数是max.poll.records和max.poll.interval.ms的关系。max.poll.records默认 500决定了单次poll()最多返回多少条消息。如果每条消息的处理时间较长比如调用第三方接口需要几百毫秒那么一次处理 500 条可能要几十秒超过max.poll.interval.ms后会被认为消费者失联触发 Rebalance。解决办法有两种调小max.poll.records比如 100 或 200或者调大max.poll.interval.ms。我更推荐前者因为它能让每轮处理时间更可控也更容易估算消费延迟。fetch.min.bytes和fetch.max.wait.ms影响消费吞吐。如果fetch.min.bytes设置成 1KBBroker 至少要攒够 1KB 数据才返回可以减少网络往返次数配合fetch.max.wait.ms500可以让消费者在数据不足时最多等待 500ms。这两个参数对消费大消息或批量消息场景帮助很大。还有一个很容易被忽视的是enable.auto.commit和auto.offset.reset。生产环境我建议关闭自动提交使用手动提交否则消息处理失败时会因为 offset 已经提交而丢消息。auto.offset.reset在消费组没有提交记录时生效如果业务不能接受从头消费要显式设置为latest。5.4 Broker 端参数num.network.threads、num.io.threads 与副本限流Broker 端调优很多人一上来就乱调线程数其实大部分场景默认值就够用。真正需要关注的是网络线程和 IO 线程的比例以及副本拉取是否成为瓶颈。num.network.threads负责处理网络请求默认是 3如果机器核数多、连接数大可以适当调大到 8 或 16。num.io.threads负责处理磁盘读写请求默认是 8一般不需要超过 CPU 核心数。不建议盲调否则线程上下文切换本身就是一种开销。副本相关的两个参数容易被忽略num.replica.fetchers和replica.fetch.max.bytes。当集群需要追赶大量副本时默认一个 follower 只有一个 fetcher 线程可能不够可以提高到 4 或 6。但要控制副本拉取带宽否则副本追赶会挤占正常的生产消费流量。我们曾经在一次重平衡中因为没有限流导致正常的业务消息延迟从 20ms 飙升到 2 秒后来通过调低replica.fetch.max.bytes才稳住。Broker 端还有一个重要参数是log.segment.bytes默认 1GB。日志段太大清理和索引重建的开销大太小会产生大量小文件。一般保持默认如果消息体大可以调大到 2GB 或 4GB。5.5 性能测试数据对比表在同样的三节点环境下我用不同参数组合做了几轮压测结果很有参考价值消息大小 1KB单 Topic 三副本配置组合Producer 吞吐平均延迟99% 延迟默认配置acks182 MB/s8 ms35 ms开启 batch64KBlinger20msacks1186 MB/s11 ms42 ms开启 batch64KBlinger20mslz4acksall172 MB/s13 ms48 ms开启 batch64KBlinger20mszstdacksall198 MB/s15 ms52 ms可以看到开启批处理和压缩后即使使用acksall吞吐依然可以远超默认配置而且可靠性更高。这组数据让我深刻体会到高可靠和性能并不是对立关系关键是找到正确的调优组合。6. Java 客户端开发与集群连接一套可靠的 Producer/Consumer 写法6.1 客户端参数与连接串的设计集群搭好了最终要通过 Java 客户端接入。官方 Java 客户端依赖坐标是dependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version3.6.2/version /dependencybootstrap.servers一定要配置多个节点地址不要只写一个。虽然一个节点也能获取到整个集群的元数据但万一这个节点就在你发送消息前宕机客户端会浪费大量时间在连接重试上。我习惯把三台 Broker 地址都写上配合合理的request.timeout.ms和delivery.timeout.ms客户端才能具备最基本的容错能力。Producer 连接串示例Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, 192.168.1.11:9092,192.168.1.12:9092,192.168.1.13:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, org.apache.kafka.common.serialization.StringSerializer); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, org.apache.kafka.common.serialization.StringSerializer); props.put(ProducerConfig.ACKS_CONFIG, all); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 65536); props.put(ProducerConfig.LINGER_MS_CONFIG, 20); props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, lz4); KafkaProducerString, String producer new KafkaProducer(props);6.2 消息可靠发送回调、重试与幂等Java Producer 的send()方法默认是异步的很多人直接producer.send(record)然后就不管了。这种方式一旦发送失败消息就静默丢失。正确做法是使用回调producer.send(record, (metadata, exception) - { if (exception ! null) { // 记录异常日志根据业务决定是否重试或落本地表 log.error(kafka send failed, exception); } });在高可靠场景下我还建议开启幂等enable.idempotencetrue。开启后Producer 每条消息会带上序列号Broker 端自动去重配合acksall可以避免重试带来的重复消息。需要注意的是幂等只保证单个 Producer 会话内的不重复跨会话仍可能出现重复消费端最好还是做幂等设计。重试参数retries默认已经是一个比较大的值Integer.MAX_VALUE不建议改小。配合delivery.timeout.ms默认 120 秒限制总发送时间避免无限重试导致消息积压在客户端内存里。如果业务对实时性要求高可以适当调小delivery.timeout.ms到 30 秒或 60 秒。6.3 消费端优雅关闭与再均衡处理消费端最常见的坑是进程被 kill 时没有关闭 Consumer导致分区没有正常释放触发不必要的 Rebalance。生产环境应该注册 JVM 关闭钩子手动调用consumer.close()。再均衡Rebalance处理也很关键。如果使用手动提交再均衡发生时可能出现部分消息已经处理完但 offset 还没提交的情况。所以ConsumerRebalanceListener里的onPartitionsRevoked回调里最好同步提交一次当前已处理的 offset避免分区转移后重复消费大量数据consumer.subscribe(Collections.singletonList(test-ha), new ConsumerRebalanceListener() { Override public void onPartitionsRevoked(CollectionTopicPartition partitions) { consumer.commitSync(); } Override public void onPartitionsAssigned(CollectionTopicPartition partitions) { // 可在这里做缓存预热等操作 } }); while (running) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { handle(record); } consumer.commitSync(); }6.4 Java 代码里的常见性能陷阱最后提几个我见过很多次的 Java 客户端性能陷阱都是真实线上踩过的。第一个陷阱是在循环里创建KafkaProducer。Producer 是重量级对象内部有缓冲区和后台发送线程反复创建会导致内存和线程资源泄漏。正确做法是使用 Spring 容器或单例模式整个应用生命周期只创建一个实例。第二个陷阱是在发送回调里做耗时操作比如写数据库、调用远程接口。回调线程是 Producer 的 IO 线程如果在里面做重操作会阻塞整个 Producer 的发送能力导致吞吐骤降。回调里只做轻量级日志和指标统计需要重操作就丢到独立线程池。第三个陷阱是消费端单线程处理所有消息。虽然 Kafka 的 Consumer 是线程不安全的但你可以根据分区数开启多个消费者实例或者使用ThreadPoolExecutor异步处理消息。理解 Kafka 的分区模型同一个分区内的消息是有序的跨分区不保证顺序。如果需要全局有序只能用一个分区加单消费者吞吐受限如果只是分区内有序多消费者并行处理是安全且高效的。第四个陷阱是没有监控消费 Lag。消费者如果处理速度跟不上生产速度Lag 会持续增长时间久了可能导致消费端内存溢出或磁盘占用过高。我建议用 Prometheus 的 kafka_exporter 或 JMX 指标定期采集消费者 Lag设置告警比如 Lag 超过 10000 就报警。我见过太多案例等发现 Kafka 有问题时Lag 已经积压到几百万恢复需要好几个小时。如果你正准备搭生产集群我建议先按第 3 章的步骤部署一套三节点环境然后原样把第 4 章的故障演练做一遍再回头按第 5 章调参。这个顺序比直接看一堆参数定义有效得多因为你已经知道每个参数在故障和压测场景下会带来什么后果。我自己就是靠这轮演练和调优把一个从能跑到能用的集群真正推到了敢接线上流量的水平。
返回列表