Kafka
基本概念
Broker消息代理 通常一个服务器就是一个broker
topic主题 相同消息存储在同一个主题中
Partition分区 同一个主题可以有一个或者多个分区 消息在写入分区的时候会分配一个序号 通常从0开始 线性增长 通常称为偏移量Offset 分区可以理解为一个有序的、不可变的提交日志
Replication副本 每个分区都有一个server作为leader 0或多个server作为follower 每个server可以作为多个分区的leader和其他若干分区的follower
Segment段
Producter生产者 生产者决定将消息发送到哪个分区
Consumer消费者 消费者通过保存offset来记录消费位置
Consumer Group消费者组 消费者组是消费者的一个属性 通过消费者组统一来点对点 发布订阅模式
下载与安装
阿里云镜像下载kafka并且安装
搭建伪分布式集群
通过在一台机器上搭建多个broker来搭建分布式集群
kafka需要依赖zookeeper 所以首先先启动zookeeper
在kafka目录下创建一个etc目录存放配置文件
将config目录下的zookeeper配置文件cp到etc目录 回到bin目录启动zookeeper
将config目录下的server.properties复制三分到etc目录下 分别修改配置文件中的broker以及端口id
使用指令创建分区 需要指定zookeeper --create用于表示创建 --topic指定topic名 并且指定副本数和副本因数
./kafka-topics.sh --zookeeper localhost:2181 --create --topic test --partitions 3 --replication-factor 2使用指令查看创建的主题
./kafka-topics.sh --zookeeper localhost:2181 --describe --topic test使用指令创建一个消费者
./kafka-console-consumer.sh --bootstrap-server localhost:9092,localhost:9095 --topic test使用指令创建一个生产者
./kafka-console-producer.sh --broker-list localhost:9092,localhost:9095 --topic test进入箭头界面就可以发送消息了
Kafka消息模型
在kafka中,分区时最小的并行单位,一个消费者可以消费多个分区,一个分区可以被多个消费者组里的消费者消费,但是,一个分区不能同时被同一个消费者组里的多个消费者消费
点对点
所有的消费者放到同一个消费者组里,这样每个分区里的消息就只能被一个消费者消费,这种模式可以实现负载均衡
发布订阅模式
每个消费者都属于不同的消费者组
分区与消息顺序
同一个生产者发送到同一分区的消息,先发送的offset肯定比后发送的offset小;同一个生产者发送到不同分区的消息,消息顺序无法保证。 消费者按照消息在分区中的顺序进行消费,kafka只保证分区内的消息顺序,不能保证分区间的消息顺序。
如果想保证所有消息的消费顺序,可以设置一个分区,但是失去了拓展性和性能;
kafka支持通过设置消息的key,相同key的消息会发送到同一个分区。
消息传递语义
最多一次————消息可能会丢失,但是永远不重复发送,consumer先提交消费位置,再读取消息,实现最多一次
最少一次————消息不会丢失,但是可能会重复,consumer先读取消息,后提交消费位置,实现最少一次
精确一次————保证消息被传递到服务端且在服务端不重复
消息传递语义需要生产者和消费者共同保证
生产者API
异步发送
使用send()进行异步发送 在虚拟机中创建一个消费者 查看消费顺序时乱序
Properties props = new Properties();
props.put("bootstrap.servers", "10.6.124.58:9092");
// props.put("acks", "all");
// props.put("retries", 0);
// props.put("linger.ms", 1);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
Producer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 20; i++)
producer.send(new ProducerRecord<String, String>("test", Integer.toString(i), Integer.toString(i)));
producer.close();同步发送
send()方法返回一个Future对象 通过调用Future对象的get方法实现阻塞 在消费者那边查看消息按照顺序消费
Properties props = new Properties();
props.put("bootstrap.servers", "10.6.124.58:9092");
// props.put("acks", "all");
// props.put("retries", 0);
// props.put("linger.ms", 1);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
Producer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 20; i++) {
Future<RecordMetadata> result = producer.send(new ProducerRecord<String, String>("test", Integer.toString(i), Integer.toString(i)));
try {
RecordMetadata recordMetadata = result.get();
} catch (InterruptedException e) {
e.printStackTrace();
} catch (ExecutionException e) {
e.printStackTrace();
}
}
producer.close();批量发送
通过设置batch.size(最大数量)与linger.ms(延迟时间)来实现批量发送 当消息满足二者任意一个时就会进行批量发送 由于设置的是同步发送 在消费者端查看到的是每秒消费一个数据
Properties props = new Properties();
props.put("bootstrap.servers", "10.6.124.58:9092");
// props.put("acks", "all");
// props.put("retries", 0);
//设置每一批消息量的大小
props.put("batch.size",16384);
//延迟时间
props.put("linger.ms", 1000);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
Producer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 20; i++) {
Future<RecordMetadata> result = producer.send(new ProducerRecord<String, String>("test", Integer.toString(i), Integer.toString(i)));
try {
RecordMetadata recordMetadata = result.get();
} catch (InterruptedException e) {
e.printStackTrace();
} catch (ExecutionException e) {
e.printStackTrace();
}
}
producer.close();acks与retries
acks就是消息的通知,当消息通过请求发送给服务端的时候,服务端需要给生产者发送一个回执,告知生产者已经收到了消息,这样生产者就知晓一个请求已经完成了。当acks为0时,生产者不会等待服务器端的任何请求,当消息被放到缓冲区时,我们就认为他成功发送了,不过可能会导致数据的丢失;当acks为1时,表示这条消息已经被服务器端的leader传输到本地了,但是不管它是否被它的follower同步,我们都通知生产者,这条消息已经被发送成功了;当acks为all时,表示不仅leader收到了这条消息,follower也已经将这条消息同步到磁盘了,这样子就保证了数据不会丢失。
retries表示当消息发送失败的时候我们可以进行重试
acks和retries配合使用可以形成不同的消息传递语义:acks=0或1是至多一次的语义;当acks=-1或all并且retries>0是至少一次的语义
消费者API
消费者与消费位置
Kafka中有一个主题 consumeroffsets 用于保存消费者消费到哪个主题、哪个分区的哪个消费位置 利于快速恢复
自动提交——至多一次
enable.auto.commit与auto.commit.interval.ms(每隔多长时间自动提交一次)用于设置自动提交
一旦消费者调用poll()方法,我们就认为消息已经被消费了,不管后面怎么处理,都认为是成功了
Properties props = new Properties();
props.setProperty("bootstrap.servers", "localhost:9092");
props.setProperty("group.id", "test");
props.setProperty("enable.auto.commit", "true");
props.setProperty("auto.commit.interval.ms", "1000");
props.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("foo", "bar"));
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());
}
手动提交——至少一次
手动提交需要设置enable.auto.commit设置为false,手动提交需要调用consumer.commitSync()方法,先处理数据然后再进行提交
Properties props = new Properties();
props.setProperty("bootstrap.servers", "localhost:9092");
props.setProperty("group.id", "test");
props.setProperty("enable.auto.commit", "false");
props.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("foo", "bar"));
final int minBatchSize = 200;
List<ConsumerRecord<String, String>> buffer = new ArrayList<>();
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
buffer.add(record);
}
if (buffer.size() >= minBatchSize) {
insertIntoDb(buffer);
consumer.commitSync();
buffer.clear();
}
}手动提交——逐条提交
使用consumer.poll(Duration.ofMillis(Long.MAX_VALUE));一次性请求尽可能多的数据
然后手动获取当前消息long lastOffset = partitionRecords.get(partitionRecords.size() - 1).offset();
然后进行逐条的提交consumer.commitSync(Collections.singletonMap(partition, new OffsetAndMetadata(lastOffset + 1)));
try {
while(running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(Long.MAX_VALUE));
for (TopicPartition partition : records.partitions()) {
List<ConsumerRecord<String, String>> partitionRecords = records.records(partition);
for (ConsumerRecord<String, String> record : partitionRecords) {
System.out.println(record.offset() + ": " + record.value());
}
long lastOffset = partitionRecords.get(partitionRecords.size() - 1).offset();
consumer.commitSync(Collections.singletonMap(partition, new OffsetAndMetadata(lastOffset + 1)));
}
}
} finally {
consumer.close();
}