Administrator
发布于 2022-07-23 / 0 阅读
0
0

kafka-202509120218

kafka-202509120218

kafka

(一)Kafka介绍

Kafka 也是是我们在开发过程中经常会使用的一种消息队列 Apache Kafka是一个分布式流处理平

台。它最初由LinkedIn开发,后来成为Apache软件基金会的一部分,并在开源社区中得到了广泛

应用。Kafka的核心概念包括Producer、Consumer、Broker、Topic、Partition和Offset。

Producer:生产者,负责将数据发送到Kafka集群。

Consumer:消费者,从Kafka集群中读取数据。

Broker:Kafka服务器实例,Kafka集群通常由多个Broker组成。

Topic:主题,数据按主题进行分类。

Partition:分区,每个主题可以有多个分区,用于实现并行处理和提高吞吐量。

Offset:偏移量,每个消息在其分区中的唯一标识。

(二)使用场景

Kafka适用于以下场景:

日志收集:集中收集系统日志和应用日志,通过Kafka传输到大数据处理系统。

消息队列:作为高吞吐量、低延迟的消息队列系统。

数据流处理:实时处理数据流,用于实时分析、监控和处理。

事件源架构:将所有的变更事件存储在Kafka中,实现事件溯源和回放。

流数据管道:构建数据管道,连接数据源和数据存储系统。

(三)Kafka的使用

前置条件

先启动zooker 服务器、再启动kafka 服务端

1.2 引入依赖

org.springframework.kafka

1.3 application.propertise配置

spring-kafka

#####################################################################

Kafka生产者配置文件

包含:必须配置、强烈建议配置和可选配置

用途:配置Kafka生产者的行为,包括消息发送方式、性能优化、可靠性保证等

注意:某些配置项的值需要根据实际生产环境进行调整

#####################################################################

###########【必须配置】###########

Kafka服务器地址,多个地址用逗号分隔(必须)

例如:localhost:9092,localhost:9093,localhost:9094

spring.kafka.producer.bootstrap-servers=localhost:9092

消息键和值的序列化器(必须)

用于将Java对象转换为Kafka中的二进制数据

spring.kafka.producer.key-

serializer=org.apache.kafka.common.serialization.StringSerializer

spring.kafka.producer.value-

serializer=org.apache.kafka.common.serialization.StringSerializer

###########【强烈建议配置】###########

可靠性配置

acks=all 所有副本都确认才算写入成功,最高可靠性

retries 发送失败时的重试次数

#enable.idempotence=true 启用幂等性,防止消息重复

#spring.kafka.producer.acks=all

#spring.kafka.producer.retries=1

#spring.kafka.producer.properties.enable.idempotence=true

性能配置

batch-size 批量发送的大小(字节)

buffer-memory 生产者缓冲区大小(字节)

compression-type 消息压缩类型(none, gzip, snappy, lz4, zstd)

spring.kafka.producer.batch-size=16384

spring.kafka.producer.buffer-memory=33554432

spring.kafka.producer.compression-type=snappy

请求配置

request.timeout.ms 等待服务器响应的最大时间

max.request.size 单个请求的最大大小

spring.kafka.producer.properties.request.timeout.ms=30000

spring.kafka.producer.properties.max.request.size=1048576

批量发送配置

linger.ms 延迟发送时间,等待更多消息一起发送

增加此值可以提高吞吐量,但会增加延迟

spring.kafka.producer.properties.linger.ms=10

发送缓冲区配置

send.buffer.bytes TCP发送缓冲区大小

receive.buffer.bytes TCP接收缓冲区大小

spring.kafka.producer.properties.send.buffer.bytes=131072

spring.kafka.producer.properties.receive.buffer.bytes=32768

###########【可选配置】###########

事务配置

如果启用事务,必须配置事务ID前缀

不同的生产者必须使用不同的事务ID

spring.kafka.producer.transaction-id-prefix=tx-1

生产者限制配置

max.block.ms 发送阻塞的最大时间

max.in.flight.requests.per.connection 单个连接最大未确认请求数

spring.kafka.producer.properties.max.block.ms=60000

#

spring.kafka.producer.properties.max.in.flight.requests.per.connection=

5

#####################################################################--

--------------------------------------------------

Kafka消费者配置文件

包含:必须配置、强烈建议配置和可选配置

用途:配置Kafka消费者的行为,包括消费方式、批量处理、会话管理等

注意:某些配置项的值需要根据实际生产环境进行调整

#####################################################################

###########【必须配置】###########

Kafka服务器地址,多个地址用逗号分隔(必须)

例如:localhost:9092,localhost:9093,localhost:9094

spring.kafka.consumer.bootstrap-servers=localhost:9092

消费者组ID,同一组的消费者协同消费消息(必须)

相同组ID的消费者消费不同分区的消息,实现负载均衡

spring.kafka.consumer.group-id=defaultConsumerGroup

消息键和值的反序列化器(必须)

用于将Kafka中的二进制数据转换为Java对象

spring.kafka.consumer.key-

deserializer=org.apache.kafka.common.serialization.StringDeserializer

spring.kafka.consumer.value-

deserializer=org.apache.kafka.common.serialization.StringDeserializer

###########【强烈建议配置】###########

消息提交方式配置

enable-auto-commit=false 关闭自动提交,防止消息丢失

ack-mode=MANUAL 手动确认模式,确保消息被正确处理后才提交

spring.kafka.consumer.enable-auto-commit=false

spring.kafka.listener.ack-mode=MANUAL

批量消费配置

type=batch 启用批量消费模式,提高消费效率

max-poll-records 每次批量消费的最大消息数

spring.kafka.listener.type=batch

spring.kafka.consumer.max-poll-records=500

会话超时配置

session.timeout.ms 消费者组会话超时时间,如果超过此时间没有心跳则认为消费者死亡

max.poll.interval.ms 两次poll之间的最大间隔,超过则认为消费者处理能力不足

spring.kafka.consumer.properties.session.timeout.ms=45000

spring.kafka.consumer.properties.max.poll.interval.ms=300000

偏移量配置

earliest: 从最早的消息开始消费

latest: 从最新的消息开始消费,保证消费者只处理最新的消息

none: 如果无偏移量则抛出异常

spring.kafka.consumer.auto-offset-reset=latest

消费者拉取配置

fetch.min.bytes 每次最小拉取大小,避免频繁拉取

fetch.max.wait.ms 当数据量不足fetch.min.bytes时,最多等待时间

spring.kafka.consumer.fetch.min.bytes=1

spring.kafka.consumer.fetch.max.wait.ms=500

心跳配置

heartbeat.interval.ms 心跳间隔时间,必须小于session.timeout.ms

spring.kafka.consumer.properties.heartbeat.interval.ms=3000

###########【可选配置】###########

并发消费配置

设置消费者线程数,提高消费能力

spring.kafka.listener.concurrency=3

消费者限制配置

max.partition.fetch.bytes 每个分区返回的最大数据量

2、简单实践

2.1生产者

2.2消费者

上面示例创建了一个生产者,发送消息到topic1,消费者监听topic1消费消息。

监听器用@KafkaListener注解,topics表示监听的topic,支持同时监听多个,用英文逗号分隔。

fetch.max.bytes 一次请求中返回的最大数据量

spring.kafka.consumer.max-partition-fetch-bytes=1048576

spring.kafka.consumer.fetch.max.bytes=52428800

@Autowired

private KafkaTemplate kafkaTemplate;

// 同步发送消息

@GetMapping("/kafka/sync")

public void sendSyncMessage(@RequestParam("message") String message) {

try {

SendResult result =

kafkaTemplate.send("topic1", message).get();

System.out.println("消息发送成功------>>>:" +

result.getRecordMetadata().offset());

} catch (Exception e) {

System.err.println("消息发送失败:" + e.getMessage());

}

}

@Component

public class KafkaConsumer {

// 消费监听

@KafkaListener(id = "tx-1",topics = {"topic1"})

public void onMessage1(ConsumerRecord record){

// 消费的哪个topic、partition的消息,打印出消息内容

System.out.println("简单消费:"+record.topic()+"-

"+record.partition()+"-"+record.value());

}

}

3.kafka事务

3.1kafka事务开启

/**

  • 配置生产者的KafkaTemplate
  • 用于发送消息到Kafka
  • */

    @Bean

    public KafkaTemplate

    kafkaTemplate(ProducerFactory producerFactory) {

    KafkaTemplate template = new KafkaTemplate<>

    (producerFactory);

    // 添加生产者监听器,用于监控消息发送状态

    template.setProducerListener(new ProducerListener

    () {

    @Override

    public void onSuccess(ProducerRecord

    producerRecord, RecordMetadata recordMetadata) {

    logger.info("消息发送成功:topic = {}, partition = {},

    offset = {}, value = {}",

    producerRecord.topic(),

    recordMetadata.partition(),

    recordMetadata.offset(),

    producerRecord.value());

    }

    @Override

    public void onError(ProducerRecord

    producerRecord, RecordMetadata recordMetadata, Exception exception) {

    logger.error("消息发送失败:topic = {}, value = {}, error =

    {}",

    producerRecord.topic(),

    producerRecord.value(),

    exception.getMessage());

    }

    });

    return template;

    }

    在Spring Kafka中,当您设置了transaction-id-prefix后,Spring会为您的KafkaTemplate配置一个

    事务性生产者。

    在 KafkaConfig 中添加事务管理器配置,并确保 ProducerFactory 支持事务

    事务消息发送

    @Bean

    public ProducerFactory producerFactory() {

    Map configProps = new HashMap<>();

    configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,

    "localhost:9092");

    configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,

    StringSerializer.class);

    configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,

    StringSerializer.class);

    configProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "tx-1");

    // 事务ID

    configProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);

    // 启用幂等性

    return new DefaultKafkaProducerFactory<>(configProps);

    }

    @Bean

    public KafkaTransactionManager

    kafkaTransactionManager() {

    return new KafkaTransactionManager<>(producerFactory());

    }

    // 发送消息(开启事务)

    @GetMapping("/kafka/transactional2")

    public void sendTransactionalMessage2(@RequestParam("message") String

    message){

    // 声明事务:出现报错的话, executeInTransaction 中包裹的 send消息是不会发

    出去

    kafkaTemplate.executeInTransaction(operations -> {

    operations.send("topic1",message + "-1");

    operations.send("topic1",message + "-2");

    3.2 kafka事务核心组件

    幂等性:保证在消息重发的时候,消费者不会重复处理。即使在消费者收到重复消息的时候,重

    复处理,也要保证最终结果的一致性。

    当Producer发送消息(x2,y2)给Broker时,Broker接收到消息并将其追加到消息流中。此时,Broke

    r返回Ack信号给Producer时,发生异常导致Producer接收Ack信号失败。对于Producer来说,会

    触发重试机制,将消息(x2,y2)再次发送,但是,由于引入了幂等性,在每条消息中附带了PID(P

    roducerID)和SequenceNumber。

    相同的PID和SequenceNumber发送给Broker,而之前Broker缓存过之前发送的相同的消息,那么

    在消息流中的消息就只有一条(x2,y2),不会出现重复发送的情况。

    事务协调器(TransactionCoordinator)

    作用:管理事务生命周期,协调事务提交或中止。

    实现机制:

    每个事务绑定一个协调器(通过事务ID哈希选择Broker)。

    维护事务状态机(TransactionState),存储在内部Topic __transaction_state。

    throw new RuntimeException("模拟事务回滚");

    });

    }

    事务日志(Transaction Log)

    作用:持久化事务状态,防止协调器宕机后数据丢失。

    存储位置:内部Topic __transaction_state,每个分区对应一个协调器。

    数据格式:事务ID、PID、状态(PrepareCommit、Completed等)、超时时间。

    事务原理流程如下:

    寻找 TC 服务地址

    Producer 会首先从 Kafka 集群中选择任意一台机器,然后向其发送请求,获取 TC 服务的地址。

    Kafka 有个特殊的事务 topic,名称为__transaction_state ,负责持久化事务消息。这个 topic 有

    多个分区,默认有50个,每个分区负责一部分事务。事务划分是根据 transaction id, 计算出该

    事务属于哪个分区。这个分区的 leader 所在的机器,负责这个事务的TC 服务地址。

    事务初始化

    Producer 在使用事务功能,必须先自定义一个唯一的 transaction id。有了 transaction id,即使

    客户端挂掉了,它重启后也能继续处理未完成的事务。

    Kafka 实现事务需要依靠幂等性,而幂等性需要指定 producer id 。所以Producer在启动事务之

    前,需要向 TC 服务申请 producer id。TC 服务在分配 producer id 后,会将它持久化到事务 topi

    c。

    发送消息

    Producer 在接收到 producer id 后,就可以正常的发送消息了。不过发送消息之前,需要先将这

    些消息的分区地址,上传到 TC 服务。TC 服务会将这些分区地址持久化到事务 topic。然后 Prod

    ucer 才会真正的发送消息,这些消息与普通消息不同,它们会有一个字段,表示自身是事务消

    息。

    发送提交请求

    Producer 发送完消息后,如果认为该事务可以提交了,就会发送提交请求到 TC 服务。Producer

    的工作至此就完成了,接下来它只需要等待响应。这里需要强调下,Producer 会在发送事务提交

    请求之前,会等待之前所有的请求都已经发送并且响应成功。

    提交请求持久化

    TC 服务收到事务提交请求后,会先将提交信息先持久化到事务 topic 。持久化成功后,服务端就

    立即发送成功响应给 Producer。然后找到该事务涉及到的所有分区,为每 个分区生成提交请

    求,存到队列里等待发送。

    发送事务结果信息给分区

    后台线程会不停的从队列里,拉取请求并且发送到分区。当一个分区收到事务结果消息后,会将

    结果保存到分区里,并且返回成功响应到 TC服务。当 TC 服务收到所有分区的成功响应后,会持

    久化一条事务完成的消息到事务 topic。至此,一个完整的事务流程就完成了。

    (四)Kafka 消息

    Producer 端丢失场景

    Producer 端为了提升发送效率,减少IO操作,发送数据的时候是将多个请求合并成一个个 Recor

    dBatch,并将其封装转换成 Request 请求「异步」将数据发送出去(也可以按时间间隔方式,达

    到时间间隔自动发送),所以 Producer 端消息丢失更多是因为消息根本就没有发送到 Kafka Brok

    er 端。

    导致 Producer 端消息没有发送成功有以下原因:

    1. 网络原因:由于网络抖动导致数据根本就没发送到 Broker 端。

    2. 数据原因:消息体太大超出 Broker 承受范围而导致 Broker 拒收消息。

    解决方法:

    更换调用方式:使用异步发送消息

    ACK 确认机制

    将 request.required.acks设置为 -1/ all

    重试次数 retries 设置为大于0的数

    重试时间 retry.backoff.ms:该参数表示消息发送超时后两次重试之间的间隔时间,避免无效的

    频繁重试,默认值为100ms, 推荐设置为300ms。

    Broker 端丢失场景

    KafkaBroker 集群接收到数据后会将数据进行持久化存储到磁盘,为了提高吞吐量和性能,采用

    的是「异步批量刷盘的策略」,也就是说按照一定的消息量和间隔时间进行刷盘。首先会将数据

    存储到 「PageCache」 中,至于什么时候将 Cache 中的数据刷盘是由「操作系统」根据自己的

    策略决定或者调用 fsync 命令进行强制刷盘,

    如果此时 Broker 宕机 Crash 掉,且选举了一个落后 Leader Partition 很多的 Follower Partition 成

    为新的 Leader Partition,那么落后的消息数据就会丢失。

    由于 Kafka 中并没有提供「同步刷盘」的方式,所以说从单个 Broker 来看还是很有可能丢失数

    据的。

    kafka 通过「多 Partition (分区)多 Replica(副本)机制」已经可以最大限度的保证数据不丢

    失,如果数据已经写入 PageCache 中但是还没来得及刷写到磁盘,此时如果所在 Broker 突然宕

    机挂掉或者停电,极端情况还是会造成数据丢失。

    解决方法:修改Broker配置

    unclean.leader.election.enable:false

    该参数表示有哪些 Follower 可以有资格被选举为 Leader , 如果一个 Follower 的数据落后 Leader

    太多,那么一旦它被选举为新的 Leader, 数据就会丢失,因此我们要将其设置为false,防止此

    类情况发生。

    replication.factor:

    该参数表示分区副本的个数。建议设置replication.factor >=3, 这样如果 Leader 副本异常 Crash

    掉,Follower 副本会被选举为新的 Leader 副本继续提供服务。

    min.insync.replicas:

    该参数表示消息至少要被写入成功到 ISR 多少个副本才算"已提交",建议设置min.insync.replica

    s > 1,这样才可以提升消息持久性,保证数据不丢失。

    另外我们还需要确保一下replication.factor > min.insync.replicas, 如果相等,只要有一个副本异

    常 Crash 掉,整个分区就无法正常工作了,因此推荐设置成:replication.factor =min.insync.repl

    icas +1, 最大限度保证系统可用性。

    Consumer 端丢失场景

    1. Consumer 拉取数据之前跟Producer 发送数据一样, 需要通过订阅关系获取到集群元数据,找

    到相关 Topic 对应的 Leader Partition 的元数据。


    评论