kafka生产者


目录
  • 原理
  • 异步发送
    • 带回调函数的异步发送流程
  • 同步发送
  • 生产者分区
    • 分区的特点
  • 分区策略
    • 1 默认的分区规则 Default'Partitioner
    • 2 自定义分区规则
  • 自定义分区器
    • 自定义分区器
    • 关联自定义分区器
  • 提高生产者吞吐量
    • 提高linger.ms的优缺点
    • 压缩的可选参数
  • Kafka数据可靠性
    • 数据的发送流程
    • ACK应答级别
      • 0:生产者发送过来的数据,不需要等数据落盘应答
      • 1:生产者发送过来的数据,Leader收到数据后应答
      • -1:生产者发送过来的数据,Leader和ISR队列里面的所有节点收齐数据后应答。
      • 真正的可靠性:
      • 可靠性总结
      • acks代码配置
  • 数据重复
    • 数据传递语义
    • 幂等性原理
      • 如何使用幂等性
    • 生产者事务
      • 事务的原理
      • 事务的API
  • 数据有序
  • 数据乱序
  • 代码

原理

异步发送

带回调函数的异步发送流程

Callback

// 2 发送数据
for (int i = 0; i < 500; i++) {
	kafkaProducer.send(new ProducerRecord<>("first", "atguigu" + i), new Callback() {
    @Override
    public void onCompletion(RecordMetadata metadata, Exception exception) {
		if (exception == null){
        	System.out.println("主题: "+metadata.topic() + " 分区: "+ metadata.partition());
       	}
    }
	});

    Thread.sleep(2);
}

同步发送

// 2 发送数据
for (int i = 0; i < 5; i++) {
	kafkaProducer.send(new ProducerRecord<>("first","atguigu"+i)).get();
}

生产者分区

kafka拦截器在生产环境中用的并不多,主要是用的Fulme的拦截器。

序列化器也不多说,主要都是字符串数据,自定义的类型较少。

分区的特点

  1. 存储的角度:便于合理使用存储资源,每个Partition在一个Broker上存储,可以把海量的数据按照分区切割成一块一块数据存储在多态Broker上。合理控制分区的任务,实现负载均衡的效果。

  2. 计算的角度:提高并行度,生产者可以以分区为单位发送数据;消费者可以以分区为单位进行消费数据。

    image-20220405155543437

分区策略

1 默认的分区规则 Default'Partitioner

2 自定义分区规则

image-20220405160104951

自定义分区器

需求:实现一个分区器实现,发送过来的数据中如果包含 atguigu,就发往 0 号分区,

不包含lsyhahaha,就发往 1 号分区

自定义分区器

import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
import java.util.Map;

public class MyPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {

        // 获取数据 atguigu  hello
        String msgValues = value.toString();

        int partition;

        if (msgValues.contains("atguigu")){
            partition = 0;
        }else {
            partition = 1;
        }

        return partition;
    }

    @Override
    public void close() {

    }

    @Override
    public void configure(Map configs) {

    }
}

关联自定义分区器

image-20220405161216477

提高生产者吞吐量

image-20220405161540207

提高linger.ms的优缺点

// 缓冲区大小
properties.put(ProducerConfig.BUFFER_MEMORY_CONFIG,33554432);
// 批次大小
properties.put(ProducerConfig.BATCH_SIZE_CONFIG,16384);
// linger.ms
properties.put(ProducerConfig.LINGER_MS_CONFIG, 1);
// 压缩
properties.put(ProducerConfig.COMPRESSION_TYPE_CONFIG,"snappy");

好处:一次拉取的数据量更大

坏处:拉去数据的延迟变高

压缩的可选参数

image-20220405162433692

Kafka数据可靠性

数据的发送流程

image-20220405162734309

ACK应答级别

0:生产者发送过来的数据,不需要等数据落盘应答

image-20220405162940228

1:生产者发送过来的数据,Leader收到数据后应答

image-20220405163022313

-1:生产者发送过来的数据,Leader和ISR队列里面的所有节点收齐数据后应答。

遇到新问题::Leader收到数据,所有Follower都开始同步数据,但有一个Follower,因为某种故障,迟迟不能与Leader进行同步,那这个问题怎么解决呢?

image-20220405163343320

真正的可靠性:

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

可靠性总结

image-20220405163524910

acks代码配置

// 设置 acks
 properties.put(ProducerConfig.ACKS_CONFIG, "1");
 // 重试次数 retries,默认是 int 最大值,2147483647
 properties.put(ProducerConfig.RETRIES_CONFIG, 3);

数据重复

数据传递语义

image-20220405163711751

image-20220405164228324

幂等性原理

概念:幂等性就是指Producer不论向Broker发送多少次重复数据,Broker端都只会持久化一条,保证了不重复。

精确一次(Exactly Once)=幂等性+至少一次(ack=-1 +分区副本数>=2 + ISR最小副本数量>=2)

image-20220405164602633

如何使用幂等性

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

生产者事务

事务的原理

事务的底层就是幂等性。所以开启事务之前,必须打开事务。

image-20220405165137901

事务的API

必须要指定事务ID(随便取,但是要保证ID全局唯一)

image-20220405165329623

数据有序

image-20220405165954944

数据乱序

解决数据乱序的问题

image-20220405170223719

只要发现数据不是有序的(如下图),就将4、5放在内存中,不会落盘到硬盘中,直到等到3也来了之后,才会将这个请求数据落盘。故无论如何,都可以保证最近5个数据都是有序的。

image-20220405170327204

代码

CustomProducerTranactions

package com.atguigu.kafka.producer;

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;

public class CustomProducerTranactions {

    public static void main(String[] args) {

        // 0 配置
        Properties properties = new Properties();

        // 连接集群 bootstrap.servers
        properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop102:9092,hadoop103:9092");

        // 指定对应的key和value的序列化类型 key.serializer
//        properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer");
        properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        // 指定事务id
        properties.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "tranactional_id_01");

        // 1 创建kafka生产者对象
        // "" hello
        KafkaProducer kafkaProducer = new KafkaProducer<>(properties);

        kafkaProducer.initTransactions();

        // 启动事务
        kafkaProducer.beginTransaction();

        try {
            // 2 发送数据
            for (int i = 0; i < 5; i++) {
                kafkaProducer.send(new ProducerRecord<>("first", "atguigu" + i));
            }

            int i = 1 / 0;

            kafkaProducer.commitTransaction();
        } catch (Exception e) {
            //
            kafkaProducer.abortTransaction();
        } finally {
            // 3 关闭资源
            kafkaProducer.close();
        }
    }
}

MyPartitioner

package com.atguigu.kafka.producer;

import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;

import java.util.Map;

public class MyPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {

        // 获取数据 atguigu  hello
        String msgValues = value.toString();

        int partition;

        if (msgValues.contains("atguigu")){
            partition = 0;
        }else {
            partition = 1;
        }

        return partition;
    }

    @Override
    public void close() {

    }

    @Override
    public void configure(Map configs) {

    }
}