ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka 基础概念速览:消息、主题、分区与消费者组

Kafka 基础概念速览:消息、主题、分区与消费者组 1. 引言Apache Kafka 是一个分布式、高吞吐、可扩展的消息队列系统广泛用于日志收集、流式处理、事件驱动架构等场景。对于初学者来说Kafka 的术语体系消息、主题、分区、消费者组等往往是最先遇到的坎。本文用通俗的语言和清晰的示例带你快速掌握 Kafka 最核心的基础概念为后续深入学习和实战打下坚实基础。2. 消息Message与批次Batch消息是 Kafka 中数据的基本单元可以理解为一条记录Record。它由三部分组成键Key可选用于决定消息被写入哪个分区相同 Key 的消息会进入同一分区从而保证顺序。值Value消息的实际内容可以是任意字节数组文本、JSON、Avro 等。时间戳Timestamp消息的创建时间用于日志保留和流处理。批次Batch是 Kafka 提高效率的关键机制。生产者不会一条一条地发送消息而是把多条消息打包成一个批次一次性发送。批次可以显著减少网络往返次数提升吞吐量但也会引入少量延迟因为要等攒够一批或达到超时时间。类比消息就像快递包裹批次就像把多个包裹装进同一辆货车再发车。3. 主题Topic与分区Partition主题Topic是消息的逻辑分类类似于数据库中的“表”。生产者把消息写入某个主题消费者从该主题读取消息。一个主题可以包含任意数量的消息。分区Partition是主题的物理分片。每个主题至少有一个分区分区内消息是有序的按写入顺序追加但分区之间不保证全局有序。分区的作用并行与扩展多个分区可以分布在集群的不同 Broker 上生产者和消费者可以并行读写从而水平扩展吞吐量。顺序保证同一分区内的消息严格有序适合需要局部有序的场景如订单状态流转。类比主题像一本杂志分区像杂志的不同分册每册内部页码连续但不同分册之间没有统一页码。下图展示了 Kafka 核心架构中生产者、Broker 集群、主题分区与消费者组之间的数据流转关系消费者组Consumer GroupBroker 集群生产者Producer主题Topic生产者 A生产者 B分区 0分区 1分区 2消费者 1消费者 2说明生产者将消息写入主题的各个分区Broker 集群负责存储与副本同步同一消费者组内的消费者各自负责不同分区实现负载均衡与水平扩展。整体形成「生产者 → 主题分区 → 消费者组」的完整数据流转链路。4. 生产者Producer与消费者Consumer生产者Producer负责把消息发布到指定的主题。生产者可以指定消息的 KeyKafka 根据 Key 的哈希值决定写入哪个分区无 Key 时采用轮询或随机策略。生产者还支持确认机制acks用于控制消息写入的可靠性级别。消费者Consumer负责从主题的分区中拉取Pull消息并处理。Kafka 采用拉取模型消费者主动向 Broker 请求数据这样消费者可以根据自身处理能力控制消费速度避免被“淹没”。下面是一个 Java 版生产者示例演示如何创建KafkaProducer、指定 Key 发送消息并设置acksall保证高可靠性importorg.apache.kafka.clients.producer.KafkaProducer;importorg.apache.kafka.clients.producer.ProducerRecord;importorg.apache.kafka.clients.producer.ProducerConfig;importorg.apache.kafka.common.serialization.StringSerializer;importjava.util.Properties;publicclassOrderProducer{publicstaticvoidmain(String[]args){// 1. 配置生产者参数PropertiespropsnewProperties();// 指定 Kafka 集群的 Broker 地址多个用逗号分隔props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,localhost:9092);// 消息 Key 的序列化器将字符串 Key 转为字节数组props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());// 消息 Value 的序列化器将消息内容转为字节数组props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());// acksall等待所有 ISR 副本确认后才视为发送成功可靠性最高props.put(ProducerConfig.ACKS_CONFIG,all);// 2. 创建 KafkaProducer 实例KafkaProducerString,StringproducernewKafkaProducer(props);try{// 3. 构造消息指定主题、Key 和 Value// 相同 Key 的消息会进入同一分区从而保证该 Key 下的消息顺序ProducerRecordString,StringrecordnewProducerRecord(order-events,order-1001,订单已创建);// 4. 发送消息异步发送返回 Futureproducer.send(record,(metadata,exception)-{if(exceptionnull){// 发送成功打印消息所在的分区和偏移量System.out.println(发送成功分区metadata.partition()偏移量metadata.offset());}else{// 发送失败打印异常信息exception.printStackTrace();}});}finally{// 5. 关闭生产者释放连接资源producer.close();}}}说明acksall配合min.insync.replicas使用可确保消息写入所有同步副本后才返回成功从而在 Broker 故障时最大限度避免数据丢失。下面是一个 Java 版消费者示例演示如何创建KafkaConsumer、订阅主题、拉取消息并手动提交偏移量importorg.apache.kafka.clients.consumer.ConsumerConfig;importorg.apache.kafka.clients.consumer.ConsumerRecord;importorg.apache.kafka.clients.consumer.ConsumerRecords;importorg.apache.kafka.clients.consumer.KafkaConsumer;importorg.apache.kafka.common.serialization.StringDeserializer;importjava.time.Duration;importjava.util.Collections;importjava.util.Properties;publicclassOrderConsumer{publicstaticvoidmain(String[]args){// 1. 配置消费者参数PropertiespropsnewProperties();// 指定 Kafka 集群的 Broker 地址与生产者保持一致props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,localhost:9092);// 消息 Key 的反序列化器将字节数组还原为字符串props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class.getName());// 消息 Value 的反序列化器将字节数组还原为字符串props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class.getName());// 消费者组 ID同一组内的消费者共同分担分区消费props.put(ConsumerConfig.GROUP_ID_CONFIG,order-group);// 关闭自动提交改为手动提交偏移量以便精确控制消费进度props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,false);// 2. 创建 KafkaConsumer 实例KafkaConsumerString,StringconsumernewKafkaConsumer(props);// 3. 订阅主题可订阅一个或多个主题consumer.subscribe(Collections.singletonList(order-events));try{// 4. 循环拉取消息并处理while(true){// 拉取消息最多阻塞 1 秒ConsumerRecordsString,Stringrecordsconsumer.poll(Duration.ofSeconds(1));for(ConsumerRecordString,Stringrecord:records){// 处理业务逻辑打印消息的 Key、分区、偏移量和内容System.out.println(收到消息Keyrecord.key()分区record.partition()偏移量record.offset()内容record.value());}// 5. 手动提交偏移量处理完本批消息后提交确保进度被保存// 同步提交会阻塞直到提交完成适合对可靠性要求较高的场景consumer.commitSync();}}finally{// 6. 关闭消费者释放连接资源consumer.close();}}}说明手动提交commitSync能精确控制偏移量提交时机避免自动提交可能带来的重复消费或消息丢失配合enable.auto.commitfalse使用适合对消费可靠性要求较高的业务场景。类比生产者是“投稿人”消费者是“读者”主题是“杂志社”分区是“分册”。5. 消费者组Consumer Group与分区再平衡Rebalance消费者组Consumer Group是 Kafka 实现“广播”和“单播”的关键。多个消费者可以组成一个组共同消费一个主题组内竞争一个分区只会被组内的一个消费者消费分区与消费者一对一映射从而实现负载均衡和水平扩展。组间广播不同消费者组可以各自独立消费同一主题的全部消息互不影响。分区再平衡Rebalance是指当消费者组成员发生变化加入、退出、崩溃或分区数量变化时Kafka 会重新分配分区与消费者的对应关系。再平衡期间消费者会短暂停止消费因此应尽量避免频繁触发。类比消费者组像一个“团队”分区像“任务”再平衡就像团队人员变动后重新分配任务。6. 副本Replica与 Leader/Follower 机制为保证高可用Kafka 为每个分区维护多个副本Replica副本分布在不同的 Broker 上。副本分为两类Leader 副本负责处理该分区的所有读写请求生产者和消费者只与 Leader 交互。Follower 副本只负责从 Leader 同步数据不对外提供服务。当 Leader 宕机时某个 Follower 会被选举为新的 Leader。ISRIn-Sync Replicas是与 Leader 保持同步的副本集合。只有 ISR 中的副本才有资格被选举为 Leader。通过副本机制Kafka 在部分 Broker 故障时仍能保证数据不丢失、服务不中断。类比Leader 像“主编”Follower 像“备份编辑”主编倒下时从备份中选一位接任。7. 偏移量Offset与提交机制偏移量Offset是分区内消息的唯一递增序号用于标识消息在分区中的位置。消费者通过记录偏移量来知道“下次该从哪里继续读”。提交机制Commit是指消费者把当前消费到的偏移量保存到 Kafka 内部主题__consumer_offsets的过程。提交方式有两种自动提交消费者定期自动提交偏移量简单但可能造成重复消费或丢失。手动提交消费者在处理完消息后手动提交可精确控制但需要处理提交失败等边界情况。偏移量管理是 Kafka 实现“至少一次”“至多一次”“精确一次”等投递语义的基础。类比偏移量像“书签”提交机制像“把书签保存到笔记本”下次翻开就能接着读。8. 总结概念一句话理解消息Kafka 中数据的基本单元批次多条消息打包发送提升吞吐主题消息的逻辑分类分区主题的物理分片支持并行与有序生产者发布消息到主题消费者从分区拉取并处理消息消费者组组内竞争、组间广播再平衡消费者变动时重新分配分区副本分区的冗余备份保证高可用Leader/Follower读写走 LeaderFollower 同步备份偏移量分区内消息的位置序号提交机制保存消费进度决定投递语义掌握这些基础概念后你就可以进一步学习 Kafka 的安装部署、生产者/消费者 API 使用、以及流处理Kafka Streams / ksqlDB等进阶内容了。希望本文能帮你迈出 Kafka 学习的第一步
返回列表