Kafka知识点总结(Java版本)
Kafka:分布式流处理平台的关键概念与技术特性
Kafka,作为一款高性能的分布式流处理平台,以其卓越的数据处理能力和灵活的架构设计,在实时数据管道和应用程序构建中扮演着重要角色。以下是Kafka的核心知识点,旨在为您揭示其背后的技术魅力。
一、基础架构与核心术语
生产者(Producer):作为消息的发起者,生产者负责将数据推送至Kafka集群,是数据流转的起点。
消费者(Consumer):作为消息的接收者,消费者从Kafka集群中提取数据,实现信息的最终消费。
代理(Broker):Kafka集群中的服务器节点,它们存储数据并处理客户端的请求,是整个系统的骨架。
主题(Topic):Kafka中消息的逻辑分类,类似于数据库中的表,用于区分不同类型的数据。
分区(Partition):主题的物理分组,每个分区都是一个有序且不可变的消息序列,增强了Kafka的并发处理能力。
偏移量(Offset):分区中消息的唯一标识,用于定位消息在分区中的位置。
二、架构优势与数据持久化
Kafka的设计理念旨在提供以下架构优势:
分布式部署:通过多Broker节点组成集群,实现数据的分布式存储与处理。
高吞吐量:利用分区机制和批量发送技术,Kafka能够实现高效的数据传输。
可扩展性:随着业务需求的增长,Kafka可通过增加Broker节点轻松扩展集群规模。
数据持久化:Kafka将消息持久化到磁盘,确保数据的长期保存,防止数据丢失。
三、副本机制与容错能力
Kafka采用副本机制,为数据提供冗余保护,提升系统的容错能力:
每个分区都有多个副本,其中一个是Leader副本,负责处理读写请求,其他为Follower副本。
当Leader副本发生故障时,Kafka能够自动进行故障转移,由其他Follower副本接替领导角色。
四、消费者组与消息消费
Kafka中的消费者组机制允许多个消费者协同工作,共同消费一个主题的消息:
每个消费者组独立消费主题中的消息,组内消费者之间不会重复消费。
通过消费者组,Kafka实现了消息的负载均衡和并行处理。
五、生产者与消费者的交互
生产者负责将消息发送到指定的主题,支持同步和异步发送模式。
消费者可以根据需要按顺序消费消息,或通过指定偏移量进行任意位置的消息消费。
六、高级特性与应用
Kafka的高级特性进一步提升了其数据处理能力:
消息压缩:支持多种压缩算法,减少网络传输数据量,提高传输效率。
消息事务:确保消息的ACID特性,满足事务性消息处理的需求。
流处理:提供流处理API,允许用户构建实时数据流处理应用,实现数据的即时转换和分析。
Kafka生产者示例代码
▼java复制代码import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; public class KafkaProducerExample { public static void main(String[] args) { // 设置Kafka生产者配置 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); // Kafka集群地址 props.put("key.serializer", StringSerializer.class.getName()); props.put("value.serializer", StringSerializer.class.getName()); // 创建Kafka生产者 KafkaProducer<String, String> producer = new KafkaProducer<>(props); // 构建消息 String topic = "test-topic"; String key = "key1"; String value = "Hello, Kafka!"; // 发送消息 ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value); producer.send(record); // 关闭生产者 producer.close(); } }
Kafka消费者示例代码
▼java复制代码import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class KafkaConsumerExample { public static void main(String[] args) { // 设置Kafka消费者配置 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); // Kafka集群地址 props.put("group.id", "test-group"); // 消费者组ID props.put("key.deserializer", StringDeserializer.class.getName()); props.put("value.deserializer", StringDeserializer.class.getName()); props.put("auto.offset.reset", "earliest"); // 从最早的offset开始消费 // 创建Kafka消费者 KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); // 订阅主题 String topic = "test-topic"; consumer.subscribe(Collections.singletonList(topic)); // 拉取消息并消费 while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, String> record : records) { System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value()); } } // 实际应用中需要正确处理消费者的关闭操作,通常在接收到关闭信号时关闭消费者 consumer.close(); } }
数据一致性
Kafka的设计目标是:高吞吐、高并发、高性能。为了做到以上三点,它必须设计成分布式的,多台机器可以同时提供读写,并且需要为数据的存储做冗余备份。
图中的主题有3个分区,每个分区有3个副本,这样数据可以冗余存储,提高了数据的可用性。并且3个副本有两种角色,Leader和Follower,Follower副本会同步Leader副本的数据。
一旦Leader副本挂了,Follower副本可以选举成为新的Leader副本, 这样就提升了分区可用性,但是相对的,在提升了分区可用性的同时,也就牺牲了数据的一致性。
我们来看这样的一个场景:一个分区有3个副本,一个Leader和两个Follower。Leader副本作为数据的读写副本,所以生产者的数据都会发送给leader副本,而两个follower副本会周期性地同步leader副本的数据,但是因为网络,资源等因素的制约,同步数据的过程是有一定延迟的,所以3个副本之间的数据可能是不同的。具体如下图所示:
此时,假设leader副本因为意外原因宕掉了,那么Kafka为了提高分区可用性,此时会选择2个follower副本中的一个作为Leader对外提供数据服务。此时我们就会发现,对于消费者而言,之前leader副本能访问的数据是D,但是重新选择leader副本后,能访问的数据就变成了C,这样消费者就会认为数据丢失了,也就是所谓的数据不一致了。
为了提升数据的一致性,Kafka引入了高水位(HW :High Watermark)机制,Kafka在不同的副本之间维护了一个水位线的机制(其实也是一个偏移量的概念),消费者只能读取到水位线以下的的数据。这就是所谓的木桶理论:木桶中容纳水的高度,只能是水桶中最短的那块木板的高度。这里将整个分区看成一个木桶,其中的数据看成水,而每一个副本就是木桶上的一块木板,那么这个分区(木桶)可以被消费者消费的数据(容纳的水)其实就是数据最少的那个副本的最后数据位置(木板高度)。
也就是说,消费者一开始在消费Leader的时候,虽然Leader副本中已经有a、b、c、d 4条数据,但是由于高水位线的限制,所以也只能消费到a、b这两条数据。
这样即使leader挂掉了,但是对于消费者来讲,消费到的数据其实还是一样的,因为它能看到的数据是一样的,也就是说,消费者不会认为数据不一致。
不过也要注意,因为follower要求和leader的日志数据严格保持一致,所以就需要根据现在Leader的数据偏移量值对其他的副本进行数据截断(truncate)操作。

