kafka生产者
目录1 默认的分区规则
提高
- 原理
- 异步发送
- 带回调函数的异步发送流程
- 同步发送
- 生产者分区
- 分区的特点
- 分区策略
- 1 默认的分区规则
Default'Partitioner - 2 自定义分区规则
- 1 默认的分区规则
- 自定义分区器
- 自定义分区器
- 关联自定义分区器
- 提高生产者吞吐量
- 提高
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的拦截器。
序列化器也不多说,主要都是字符串数据,自定义的类型较少。
分区的特点
-
存储的角度:便于合理使用存储资源,每个Partition在一个Broker上存储,可以把海量的数据按照分区切割成一块一块数据存储在多态Broker上。合理控制分区的任务,实现负载均衡的效果。
-
计算的角度:提高并行度,生产者可以以分区为单位发送数据;消费者可以以分区为单位进行消费数据。

分区策略
1 默认的分区规则 Default'Partitioner
2 自定义分区规则

自定义分区器
需求:实现一个分区器实现,发送过来的数据中如果包含 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) {
}
}
关联自定义分区器

提高生产者吞吐量

提高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");
好处:一次拉取的数据量更大
坏处:拉去数据的延迟变高
压缩的可选参数

Kafka数据可靠性
数据的发送流程

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

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

-1:生产者发送过来的数据,Leader和ISR队列里面的所有节点收齐数据后应答。
遇到新问题::Leader收到数据,所有Follower都开始同步数据,但有一个Follower,因为某种故障,迟迟不能与Leader进行同步,那这个问题怎么解决呢?

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

acks代码配置
// 设置 acks
properties.put(ProducerConfig.ACKS_CONFIG, "1");
// 重试次数 retries,默认是 int 最大值,2147483647
properties.put(ProducerConfig.RETRIES_CONFIG, 3);
数据重复
数据传递语义


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

如何使用幂等性
开启参数 enable.idempotence 默认为 true,false 关闭。
生产者事务
事务的原理
事务的底层就是幂等性。所以开启事务之前,必须打开事务。

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

数据有序

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

代码
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) {
}
}