• Kafka (四) --------- 生产经验



    一、生产者如何提高吞吐量

    在这里插入图片描述

    package com.fancy.kafka.producer;
    import org.apache.kafka.clients.producer.KafkaProducer;
    import org.apache.kafka.clients.producer.ProducerRecord;
    import java.util.Properties;
    public class CustomProducerParameters {
    	public static void main(String[] args) throws InterruptedException {
    		// 1. 创建 kafka 生产者的配置对象
    		Properties properties = new Properties();
    		// 2. 给 kafka 配置对象添加配置信息:bootstrap.servers
    		properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop102:9092");
    		// key,value 序列化(必须):key.serializer,value.serializer
    		properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
    		properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
    		// batch.size:批次大小,默认 16K
    		properties.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
    		// linger.ms:等待时间,默认 0
    		properties.put(ProducerConfig.LINGER_MS_CONFIG, 1);
    		// RecordAccumulator:缓冲区大小,默认 32M:buffer.memory
    		properties.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
    		// compression.type:压缩,默认 none,可配置值 gzip、snappy、lz4 和 zstd
    		properties.put(ProducerConfig.COMPRESSION_TYPE_CONFIG,"snappy");
    		// 3. 创建 kafka 生产者对象
    		KafkaProducer<String, String> kafkaProducer = new KafkaProducer<String, String>(properties);
    		// 4. 调用 send 方法,发送消息
    		for (int i = 0; i < 5; i++) {
    			kafkaProducer.send(new ProducerRecord<>("first","fancyry " + i));
    		}
    		// 5. 关闭资源
    		kafkaProducer.close();
    	}
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31

    测试:
    ①在 hadoop102 上开启 Kafka 消费者。

    [fancyry@hadoop103 kafka]$ bin/kafka-console-consumer.sh --bootstrap-server hadoop102:9092 --topic first
    
    • 1

    ②在 IDEA 中执行代码,观察 hadoop102 控制台中是否接收到消息。

    [fancyry@hadoop102 kafka]$ bin/kafka-console-consumer.sh --bootstrap-server hadoop102:9092 --topic first
    fancyry 0
    fancyry 1
    fancyry 2
    fancyry 3
    fancyry 4
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6

    二、数据可靠性

    ACK 应答级别

    在这里插入图片描述
    思考:Leader收到数据,所有Follower都开始同步数据,但有一个Follower,因为某种故障,迟迟不能与Leader进行同步,那这个问题怎么解决呢?

    Leader 维护了一个动态的 in-sync replica set(ISR),意为和 Leader 保持同步的 Follower+Leader集合( leader:0,isr:0, 1, 2)。如果 Follower 长时间未向 Leader 发送通信请求或同步数据,则
    该 Follower 将被踢出 ISR。该时间阈值由 replica.lag.time.max.ms 参数设定,默认30s。例如2超时,(leader:0, isr:0,1)。这样就不用等长期联系不上或者已经故障的节点。

    数据可靠性分析

    如果分区副本设置为 1个 ,或 者ISR里应答的最小副本数量( min.insync.replicas 默认为1)设置为1,和ack=1的效果是一样的,仍然有丢数的风险(leader:0,isr:0)。

    数据完全可靠条件 = ACK级别设置为-1 + 分区副本大于等于2 + ISR里应答的最小副本数量大于等于2
    
    • 1

    可靠性总结

    • acks=0,生产者发送过来数据就不管了,可靠性差,效率高;
    • acks=1,生产者发送过来数据Leader应答,可靠性中等,效率中等;
    • acks=-1,生产者发送过来数据Leader和ISR队列里面所有Follwer应答,可靠性高,效率低;

    在生产环境中,acks=0很少使用;acks=1,一般用于传输普通日志,允许丢个别数据;acks=-1,一般用于传输和钱相关的数据,对可靠性要求比较高的场景。

    在这里插入图片描述

    package com.fancy.kafka.producer;
    import org.apache.kafka.clients.producer.KafkaProducer;
    import org.apache.kafka.clients.producer.ProducerRecord;
    import java.util.Properties;
    public class CustomProducerAck {
    	public static void main(String[] args) throws InterruptedException {
    		// 1. 创建 kafka 生产者的配置对象
    		Properties properties = new Properties();
    		// 2. 给 kafka 配置对象添加配置信息:bootstrap.servers
    		properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop102:9092");
    		// key,value 序列化(必须):key.serializer,value.serializer
    		properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    		properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    		// 设置 acks
    		properties.put(ProducerConfig.ACKS_CONFIG, "all");
    		// 重试次数 retries,默认是 int 最大值,2147483647
    		properties.put(ProducerConfig.RETRIES_CONFIG, 3);
    		// 3. 创建 kafka 生产者对象
    		KafkaProducer<String, String> kafkaProducer = new
    		KafkaProducer<String, String>(properties);
    		// 4. 调用 send 方法,发送消息
    		for (int i = 0; i < 5; i++) {
    			kafkaProducer.send(new
    			ProducerRecord<>("first","atguigu " + i));
    		}
    		// 5. 关闭资源
    		kafkaProducer.close();
    	}
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29

    三、数据去重

    1. 数据传递语义

    至少一次(At Least Once)= ACK级别设置为-1 + 分区副本大于等于2 + ISR里应答的最小副本数量大于等于2

    最多一次(At Most Once)= ACK级别设置为0

    总结:

    At Least Once 可以保证数据不丢失,但是不能保证数据不重复;
    At Most Once 可以保证数据不重复,但是不能保证数据不丢失。

    精确一次 (Exactly Once) :对于一些非常重要的信息,比如和钱相关的数据,要求数据既不能重复也不丢失。

    Kafka 0.11版本以后,引入了一项重大特性:幂等性和事务。

    2. 幂等性

    幂等性就是指Producer不论向Broker发送多少次重复数据,Broker端都只会持久化一条,保证了不重复。精确一次(Exactly Once) = 幂等性 + 至少一次( ack=-1 + 分区副本数>=2 + ISR最小副本数量>=2)

    重复数据的判断标准:具有相同主键的消息提交时,Broker只会持久化一条。其中PID是Kafka每次重启都会分配一个新的;Partition 表示分区号;Sequence Number是单调自增的。

    所以幂等性只能保证的是在单分区单会话内不重复

    在这里插入图片描述
    如何使用幂等性 ?

    开启参数 enable.idempotence 默认为 true,false 关闭。

    3. 生产者事务

    Kafka 事务原理

    说明:开启事务,必须开启幂等性。

    在这里插入图片描述

    Kafka 的事务一共有如下 5 个 API

    // 1 初始化事务
    void initTransactions();
    // 2 开启事务
    void beginTransaction() throws ProducerFencedException;
    // 3 在事务内提交已经消费的偏移量(主要用于消费者)
    void sendOffsetsToTransaction(Map<TopicPartition, OffsetAndMetadata> offsets, String consumerGroupId) throws ProducerFencedException;
    // 4 提交事务
    void commitTransaction() throws ProducerFencedException;
    // 5 放弃事务(类似于回滚事务的操作)
    void abortTransaction() throws ProducerFencedException;
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10

    单个 Producer,使用事务保证消息的仅一次发送

    package com.fancy.kafka.producer;
    import org.apache.kafka.clients.producer.KafkaProducer;
    import org.apache.kafka.clients.producer.ProducerRecord;
    import java.util.Properties;
    public class CustomProducerTransactions {
    	public static void main(String[] args) throws InterruptedException {
    		// 1. 创建 kafka 生产者的配置对象
    		Properties properties = new Properties();
    		// 2. 给 kafka 配置对象添加配置信息
    		properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"hadoop102:9092");
    		// key,value 序列化
    		properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    		properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    		// 设置事务 id(必须),事务 id 任意起名
    		properties.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "transaction_id_0");
    		// 3. 创建 kafka 生产者对象
    		KafkaProducer<String, String> kafkaProducer = new KafkaProducer<String, String>(properties);
    		// 初始化事务
    		kafkaProducer.initTransactions();
    		// 开启事务
    		kafkaProducer.beginTransaction();
    		try {
    			// 4. 调用 send 方法,发送消息
    			for (int i = 0; i < 5; i++) {
    				// 发送消息
    				kafkaProducer.send(new ProducerRecord<>("first", "fancy " + i));
    			}
    			// int i = 1 / 0;
    			// 提交事务
    			kafkaProducer.commitTransaction();
    		} catch (Exception e) {
    			// 终止事务
    			kafkaProducer.abortTransaction();
    		} finally {
    			// 5. 关闭资源
    			kafkaProducer.close();
    		}
    	}
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39

    四、数据有序

    在这里插入图片描述

    五、数据乱序

    kafka在1.x版本之前保证数据单分区有序,条件如下:

    max.in.flight.requests.per.connection=1(不需要考虑是否开启幂等性)。
    
    • 1

    kafka在1.x及以后版本保证数据单分区有序,条件如下:

    A、未开启幂等性

    max.in.flight.requests.per.connection需要设置为1
    • 1

    B、开启幂等性

    max.in.flight.requests.per.connection需要设置小于等于5
    • 1

    原因说明:因为在 kafka1.x 以后,启用幂等后,kafka 服务端会缓存 producer 发来的最近 5 个request 的元数据,故无论如何,都可以保证最近 5 个request的数据都是有序的。

    在这里插入图片描述

  • 相关阅读:
    #{}和${}的区别
    【Vue基础-数字大屏】自定义主题
    经典文献阅读之--EGO-Planner(无ESDF的四旋翼局部规划器)
    【分布式应用】消息队列之卡夫卡 + EFLFK集群部署
    [Spring Boot] 集成Nacos
    tp6+vue-elementui-admin实现前后端权限分离框架
    树表的查找
    快速上手Linux核心命令(九):文件备份与压缩
    ROS2——初识ROS2(一)
    算法工程师老潘总结的一些经验
  • 原文地址:https://blog.csdn.net/m0_51111980/article/details/126318910