3 Kafka的重要概念以及策略

发布时间:2026/7/26 19:06:49
3 Kafka的重要概念以及策略 引言在设计一个消息引擎系统时需要考虑的两个重要因素1.消息设计2.传输协议消息设计1.在设计消息时一定要考虑语义的清晰和格式上的通用性2.一条消息要有能够完整清晰表达业务的能力不能是含糊不清、语义不明 甚至无法处理的3.同时为了更好地表达语义以及最大限度地提高重用性消息通常都采用结构化的方式进行设计消息的结构化比如SOAP协议中的消息就采用了XML格式而Web Service也支持JSON格式的消息Kafka的消息是用二进制方式来保存的但依然是结构化的消息两种消息引擎消息队列模型发布/订阅模型消息队列模型是基于队列提供消息传输服务的多用于进程间通信以及线程间通信。该模型定义了消息队列(queue)、发送者(sender)和接收者(receiver)提供了一种点对点(point-to-point,p2p)的消息传递方式即发送者发送每条消息到队列的指定位置接收者从指定位置获取消息。一旦消息被消费(consumed)就会从队列中移除该消息。每条消息由一个发送者生产出来且只被一个消费者(consumer)处理——发送者和消费者之间是一对一的关系。发布/订阅模型它有主题topic的概念一个 topic 可以理解为逻辑语义相近的消息的容器。发布者将消息生产出来发送到指定的 topic 中所有订阅了该 topic 的订阅者都可以接收到该topic下的所有消息。如何做到高吞吐量、低延时的呢1.首先Kafka的写入操作是很快的这主要得益于它对磁盘的使用方法的不同。2.虽然Kafka会持久化所有数据到磁盘但本质上每次写入操作其实都只是把数据写入到操作系统的页缓存page cache中然后由操作系统自行决定什么时候把页缓存中的数据写回磁盘上。该设计有3个主要优势1.操作系统页缓存是在内存中分配的所以消息写入的速度非常快2.Kafka不必直接与底层的文件系统打交道。所有烦琐的I/O操作都交由操 作系统来处理3.Kafka写入操作采用追加写入append的方式避免了磁盘随机写操 作。对于普通的物理磁盘非固态硬盘而言我们总是认为磁盘的读/写操作是很慢的。事实上普通SAS磁盘随机读/写的吞吐量的确是很慢的但是磁盘的顺序读/写操作其实是非常快的它的速度甚至可以匹敌内存的随机 I/O 速度.鉴于这一事实Kafka 在设计时采用了追加写入消息的方式即只能在日志文件末尾追加写入新的消息且不允许修改已写入的消息因此它属于典型的磁盘顺序访问型操作所以Kafka 消息发送的吞吐量是很高的。实际使用过程中可以很轻松地做到每秒写入几万甚至几十万条消息消费端如何做到高吞吐量、低延时?Kafka是把消息写入操作系统的页缓存中的。那么同样地Kafka 在读取消息时会首先尝试从 OS的页缓存中读取如果命中便把消息经页缓存直接发送到网络的 Socket上。这个过程就是利用 Linux 平台的 sendfile 系统调用做到的而这种技术就是大名鼎鼎的零拷贝Zero Copy技术.除了零拷贝技术Kafka 由于大量使用页缓存故读取消息时大部分消息很有可能依然保存在页缓存中因此可以直接命中缓存不用“穿透”到底层的物理磁盘上获取消息从而极大地提升了消息读取的吞吐量。事实上如果我们监控一个经过良好调优的 Kafka生产集群便可以发现即使是那些有负载的 Kafka服务器其磁盘的读操作也很少这是因为大部分的消息读取操作会直接命中页缓存。kafka的故障转移Kafka 服务器支持故障转移的方式就是使用会话机制。每台 Kafka 服务器启动后会以会话的形式把自己注册到ZooKeeper 服务器上。一旦该服务器运转出现问题与ZooKeeper 的会话便不能维持从而超时失效此时 Kafka集群会选举出另一台服务器来完全代替这台服务器继续提供服务。服务伸缩性所谓伸缩性英文名是 scalability。根据 Java 大神 Brian Goetz 在其经典著作 Java Concurrency in Practice 中的定义伸缩性表示向分布式系统中增加额外的计算资源比如CPU、内存、存储或带宽时吞吐量提升的能力。举一个例子来说对于计算密集型computation-intensive的业务而言CPU的消耗一定是最大的这类系统上的操作我们称之为 CPU-bound。那么如果一个 CPU 的运算能力是 U我们自然希望两个 CPU 的运算能力是2U即可以线性地扩容计算能力这种线性伸缩性是最理想的状态但在实际中几乎不可能达到毕竟分布式系统中有很多隐藏的“单点”瓶颈制约了这种线性的计算能力扩容。**阻碍线性扩容的一个很常见的因素就是状态的保存。**不论是哪类分布式系统集群中的每台服务器一定会维护很多内部状态。如果由服务器自己来保存这些状态信息则必须要处理一致性的问题。相反如果服务器是无状态的状态的保存和管理交于专门的协调服务来做比如ZooKeeper那么整个集群的服务器之间就无须繁重的状态共享这极大地降低了维护复杂度。倘若要扩容集群节点只需简单地启动新的节点机器进行自动负载均衡就可以了。Kafka正是采用了这样的思想——每台Kafka服务器上的状态统一交由ZooKeeper保管。扩展 Kafka 集群也只需要一步启动新的 Kafka 服务器即可。当然这里需要言明的是在Kafka 服务器上并不是所有状态都不保存它只保存了很轻量级的内部状态因此在整个集群间维护状态一致性的代价是很低的。1.结论1.zookeeper在kafka集群中的充当的角色是一个协调者的角色存储了Kafka的元数据2.一个物理机机器上如果只启动了一个kafka进程那么就可以认为一个broker就代表一个kafka节点3.broker是有角色的角色名称是controller 多个broker之间如何协作以及创建topic,partition如何管理等是由controller去协调的4.如何在多个broker中选举一个controller,是由zookeeper决定的原理是多个broker去zookeeper中 注册争夺锁资源谁获取了锁资源谁就是controller2.broker和partition介绍1.partition和consumer的关系可以是11的关系也可以是n(partition):1(consumer)partition:consumerN:1的是绝对不允许的。 如果多个consumer,超过分区数量的时候多余的consumer会是空闲状态3.kafka 消费者不同的组之间可以重复消费同一个分区内的同一条消息么可以不同消费组Group ID 彼此完全隔离不同组的消费者能重复消费同一个分区、同一条消息同一组内 不会重复消费。3.1 核心规则消费组隔离机制Kafka 消息的消费隔离粒度是消费组1.每条消息会被所有不同消费组各自完整消费一遍2.同一个消费组内 一个分区同一时刻只会分配给组内一个消费者 组内多个消费者分摊分区组内不会重复消费场景演示假设Topictest_topic1个分区 消息msg0、msg1、msg2场景 1两个不同消费组消费组 group-A → 消费者 A1 消费组 group-B → 消费者 B1现象A1 和 B1 都会完整收到 msg0、msg1、msg2两条流互不干扰、各自维护自己的 offset。group-A 的 offset 存在 __consumer_offsets 中只作用于本组group-B 有独立的 offset 记录和 A 组无关场景 2同一个消费组两个消费者组 group-A 下A1、A2Topic 只有 1 个分区 → 分区只会分配给其中一个消费者另一消费者收不到消息自然不会重复消费。4. 如何保证同一个消费组内的消息不被重复消费前提Kafka 同一消费组天生设计就是「一条消息只被组内一个消费者消费」正常流程不会重复消费出现重复消费基本是Rebalance、offset 提交异常、客户端异常宕机导致。4.1 组内不重复的原生机制分区分配规则一个分区在同一时间只会被消费组内一个消费者持有组内其他消费者无法消费该分区从架构上避免并行重复消费。位移 (offset) 机制每个「消费组 Topic 分区」维护独立 offset标记当前已消费到的位置消费者重启 / 重分配后从已提交 offset 继续消费。组内重复消费的三大根因1.消费者宕机、进程被杀、网络超时 → 未正常提交 offset2.两次 poll 间隔过长 → 触发 Rebalance分区被重新分配3.offset 提交超时/失败但业务已处理完消息4.2 核心解决方案优先级从高到低正确配置消费参数基础防线必配1.关闭自动提交使用手动提交生产第一准则 enable.auto.committrue 是重复消费高发区自动提交是定时后台提交可能出现业务已处理消息 → 还未到自动提交时间 → 消费者宕机/Rebalance → 重启重消费推荐配置# 关闭自动提交核心 enable.auto.commitfalse # 会话超时检测消费者是否存活 session.timeout.ms10000# 心跳发送间隔建议为 session.timeout.ms 的1/3heartbeat.interval.ms3000# 两次 poll 最大间隔处理消息必须在该时间内完成否则被踢出组 max.poll.interval.ms300000# 单次 poll 拉取消息数控制批量大小避免批次过大处理超时 max.poll.records200参数解读 避坑规范 offset 提交时机代码层核心原则消息业务处理完全成功后再提交 offset处理失败绝不提交。批量同步提交中小流量、可靠性优先*//commitSync() 同步提交提交失败会抛异常可重试适合对一致性要求高的业务while(true){ConsumerRecordsString,Stringrecordsconsumer.poll(Duration.ofMillis(100));if(records.isEmpty())continue;// 1. 逐批处理所有消息for(ConsumerRecordString,Stringrecord:records){handleBusiness(record);// 执行业务逻辑}// 2. 整批处理成功再提交 offsettry{consumer.commitSync();}catch(RetriableExceptione){// 可重试异常重试提交consumer.commitSync();}catch(Exceptione){// 提交失败记录日志、告警人工介入e.printStackTrace();}}批量异步提交高吞吐场景commitAsync() 不阻塞线程吞吐更高无法即时感知提交结果建议搭配回调记录异常consumer.commitAsync((offsets,exception)-{if(exception!null){// 提交失败告警、落日志System.err.println(offset 提交失败exception.getMessage());}});规避 Rebalance重复消费重灾区Rebalance分区重分配是组内重复消费最常见诱因尽量减少 Rebalance 发生。# 优雅停止推荐 kill-15进程ID# 禁止强制杀死丢失未提交offset kill-9进程ID兜底方案消息幂等最终保障必做5 offsetoffset首先需要进行持久化老版本中会维护在zk中zk还是种协调轻业务不应该涉及业务的存储offset就是跟业务相关的数据因此在后续的kafka中在broker中维护了另外的一个topic,叫做consumer offset默认是50个分区自己维护offset. 是基于组来维护的offset可以存第三方如redis,mysql等