liuyulin
发布于 2024-02-02 / 41 阅读
0
0

Kafka

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();
     }
​


评论