在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述

重要概念:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

应用场景——面试:

好处1-异步处理,在调用需要一定时间相应的服务之前,先通过消息队列返回给用户一个信息–更快:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

好处2–功能解耦–提高稳定性–一个服务调用其他服务的时候,使用MQ向需要调用的微服务发送信息来请求服务,这样使得稳定性更高00不会使某个微服务失效导致整个功能都失效:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

kafka优势==大访问量–kafka来接受高访问避免大量请求直接进入数据库倒是数据库崩溃:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

大数据采集–直接把用户的操作的日志信息通过消息队列发送给服务器–然后服务器后台再根据主题主动拉取消息—来分析

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

kafka4.0—启动。修改配置:

4.0完全丢弃了zookeeper的依赖(不过是用zookeeper也能用,现在zookeeper相关的配置文件与KRaft模式的配置文件好像是共用的,本人菜鸟没试过),本文是在完全的window系统中安装的,如果有linux环境或者主机上有docker也可以用,官网上有完整的教程,但是没有windows的–搞了好久才搞好。哭

官网网址:

下载:    https://kafka.apache.org/downloads
文档:    https://kafka.apache.org/documentation/#quickstart

最好下载二进制版本-能够快速上手:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

保证config\server.properties中:

process.roles=controller,broker
node.id=1
listeners=PLAINTEXT://:9092,CONTROLLER://:9093
controller.listener.names=CONTROLLER  
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT  
controller.quorum.voters=1@localhost:9093
log.dirs=D:\\kafka_4.0\\kafka_2.13-4.0.0\\logs
num.partitions=1        

保证这个路径存在,不存在自己创建:

log.dirs=D:\\kafka_4.0\\kafka_2.13-4.0.0\\logs

Kafka 4.0 安装配置全流程总结(Windows 环境)

1. 初始准备
  • 下载二进制包
    确认下载的是官方二进制包(如 kafka_2.13-4.0.0.tgz),而非源码包。

  • 安装 JDK 17
    卸载 JDK 20,安装 JDK 17 并配置环境变量:

    cmd

    复制

    set JAVA_HOME=D:\JDKS\JDK17
    set PATH=%JAVA_HOME%\bin;%PATH%
    
2. 解压与目录结构
  • 解压到 D:\kafka_4.0\kafka_2.13-4.0.0,确保路径无空格或中文。
  • 关键目录:
    • bin\windows:脚本文件
    • config:配置文件
    • logs(需手动创建):数据存储目录
3. 必要配置修改
(1) 修改 config/server.properties!!最好把文件中的中文注解删除,可能会有编码错误

properties

复制

# KRaft 模式核心配置---声明当前节点既是controller又是普通broker
process.roles=controller,broker
node.id=1
listeners=PLAINTEXT://:9092,CONTROLLER://:9093
controller.listener.names=CONTROLLER
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
controller.quorum.voters=1@localhost:9093
log.dirs=D:\\kafka_4.0\\kafka_2.13-4.0.0\\logs  # 使用双反斜杠或正斜杠
num.partitions=1
(2) 配置 Log4j 2.x
  • 复制模板并重命名:(这一步不需要,亲测两个yaml都可以)

    cmd

    复制

    copy config\tool-log4j2.yaml config\log4j2.yaml
    
  • 删除旧版 Log4j 1.x 配置(如果config目录下由 log4j.properties这种文件就删除,没有则不用管)。

4. 初始化存储目录

cmd命令行中

先进入你的bin\window目录,然后启动命令行

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

cd D:\kafka_4.0\kafka_2.13-4.0.0\bin\windows

kafka-storage.bat random-uuid  # 生成集群ID,可能失败输出一个异常,如果成功则会输出一个UUID,则在下一条命令-t中写入生成的id(如 ABC123...)

kafka-storage.bat format -t ABC123... -c ..\..\config\server.properties    #初始化存储空间,若上一条失败则手动赋予一个集群id,ABC123...;-c指定server文件的地址,集群部署时每个节点的集群id字段必须相同,执行完成后效果创建一个broker
  • 生成uuid
  • 这个异常无所谓–正常生成uuid=pqGyUGPxSxarYvYuBam9ww外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传
  • 存储:记得使用刚刚生成的uuid
  • 外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传
  • 验证:检查 D:\kafka_4.0\kafka_2.13-4.0.0\logs 是否生成 meta.properties–存放uuid-即集群的id 和 __cluster_metadata-0–一个broker容器。
  • cluster.id即为集群id:
  • 外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传
  • 这一步很可能报错:但没影响-这是long4j日志报的错:warning。。。什么的,只要检查你的logs目录下有没有生成新的东西就行

几个重要的id:

cluster.id —集群id,每一个卡夫卡集群下的每个broker都要保证同一个集群id

directory.id --Kafka 节点的目录的唯一标识符。• 作用:用于标识 Kafka 节点的存储目录。在 Kafka 中,每个节点都有一个独立的存储目录,用于存储日志文件、元数据等信息。

node.id ----节点的额唯一标识,每个节点都有一个唯一的 node.id

5. 启动 Kafka

cmd命令:

第二个参数:“server.properties的路径务必注意,…\表示上一级。请根据你的目录具体更改”

kafka-server-start.bat ..\..\config\server.properties
  • 成功标志:日志输出 [KafkaRaftServer] started 且无致命错误。:
  • 这种错误完全正常:
  • 外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传
  • 正确结果:日志输出没有ERROR且打印了好几行,最后一行:外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

检查你的logs目录下有没有这个文件:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

有了就成功了

6. 测试功能:------在kafka中,一条记录或消息称为“事件(Event)”

新开一个cmd,还是进入windows目录

# 创建 Topic主题,主题名:test,必填项-指定Kafka服务的端口号
kafka-topics.bat --create --topic test --bootstrap-server localhost:9092 

#查看所有主题
kafka-topics.bat --bootstrap-server localhost:9092 --list

#查看test详情:
kafka-topics.bat --bootstrap-server localhost:9092 --topic test --describe

# 生产消息
kafka-console-producer.bat --bootstrap-server localhost:9092 --topic test

# 消费消息
kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic test --from-beginning
7. 关机与重启
  • 关机:按 Ctrl+C 停止 Kafka。
  • 重启:直接运行 kafka-server-start.bat,无需再次初始化。
  • 持久化:确保不删除 logs 目录下的文件。

关键注意事项

  1. 路径一致性:所有配置中的路径使用绝对路径,避免混用 /\
  2. 端口冲突:确保 9092(Broker)和 9093(Controller)未被占用。
  3. 权限问题:赋予 logs 目录完全控制权限。
  4. 日志警告
    • Reconfiguration failed:Log4j 内部警告,可忽略。
    • DEPRECATED: Log4j 1.x:已迁移到 Log4j 2.x,无需处理。

常见问题速查

问题现象 解决方案
启动时报 meta.properties not found 重新执行 kafka-storage.bat format
端口冲突 修改 server.properties 中的 listeners 端口
Java 版本错误 检查 java -version 是否为 JDK 17
日志目录未生成 检查 log.dirs 路径权限和拼写

按此流程操作后,Kafka 将保持持久化状态,重启后无需重复配置。

下一次重启:只需要输入:

kafka-server-start.bat ..\..\config\server.properties

一些命令:

主题详情:

kafka-topics.bat --bootstrap-server localhost:9092 --topic test --describe

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

复制主题(副本):到节点2<新的分区号>

kafka-topics.bat --bootstrap-server localhost:9092 --topic test --alter --partitions 2

window删除主题命令会报错,还是需要在linux中部署

生产数据(发送消息)

# 生产消息
kafka-console-producer.bat --bootstrap-server localhost:9092 --topic test

# 消费消息(到指定的主题中-接收消息)

# 消费消息---from-beginning获取队列中的最早消息开始读
kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic test --from-beginning

一个生产一个接受:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

特殊–需要先创建若干 __consumer_offsets主题(负责指定偏移量,转发–不知道其他的怎么样):

kafka-topics.bat --bootstrap-server localhost:9092 --describe --topic __consumer_offsets

成功接受所有消息:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

在springboot中整合使用:

导入依赖:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

配置文件中添加:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

创建生产者:

public class FirstkafkaProducter {
    public static void main(String[] args){
        // TODO 配置属性集合,用于存放Kafka生产者(消息发送者)的配置信息---连接Kafka集群9092,将消息的Key和Value序列化,
        Map<String, Object> configMap = new HashMap<>();
        // TODO 配置属性:Kafka服务器集群地址
        configMap.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        // TODO 配置属性:Kafka生产的数据为键值对,消息通过网络传输,因此需要先将消息序列化,所以在生产数据进行传输前需要分别对K,V进行对应的序列化操作
//        kafka提供自带的序列化类:"org.apache.kafka.common.serialization.StringSerializer"直接传入即可
        //对key进行序列化,常量ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG代表key的序列化类
        configMap.put(
                ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
                "org.apache.kafka.common.serialization.StringSerializer");
        configMap.put(
                ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
                "org.apache.kafka.common.serialization.StringSerializer");

        // TODO 根据配置属性集合configMap,创建Kafka生产者对象(消息发送者对象producer),建立Kafka连接,若指定的主题不存在,则会先创建
        //      构造对象时,需要传递配置参数,一帮封装到集合中,configMap
        KafkaProducer<String, String> producer = new KafkaProducer<>(configMap);
        // TODO 准备数据(消息),定义泛型
        //      构造对象时需要传递 【Topic主题名称】,【Key】,【Value--真正的的消息】三个参数
        for (int i = 0; i < 10; i++) {
            ProducerRecord<String, String> record = new ProducerRecord<String, String>(
                    "test", "key" + i, "value" + i
            );
            // TODO 生产(发送)数据
            producer.send(record);
        }

        // TODO 关闭生产者连接
        producer.close();
    }

}

消费者:

public class FirstkafkaConsumer {
    public static void main(String[] args) {
        // TODO 配置属性集合---连接Kafka集群9092,将消息的Key和Value反序列化
        Map<String, Object> configMap = new HashMap<String, Object>();

        // TODO 配置属性:Kafka集群地址
        configMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        // TODO 配置属性: Kafka传输的数据为KV对,所以需要对获取的数据键和值分别进行反序列化
        configMap.put(
                ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
                "org.apache.kafka.common.serialization.StringDeserializer");
        configMap.put(
                ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
                "org.apache.kafka.common.serialization.StringDeserializer");

        // TODO 配置属性: 新!!读取数据的位置(偏移量) ,取值为earliest(最早),latest(最晚);生产者生产消息后,分区中的索引位指向当前消息的后一位(空),而消费者初次消费时,消费者读取时的offset偏移量就是当前的索引位,因此为空读不到。这个设置后,将消费者读取时的offset重置为第一位索引位
        configMap.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"earliest");

        // TODO 配置属性: 消费者组
        configMap.put("group.id", "wangjj");
        // TODO 配置属性: 自动提交偏移量
        configMap.put("enable.auto.commit", "true");
        
        //获取消费者对象
        KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(configMap);

        // TODO 消费者订阅指定接收哪些主题的数据
        consumer.subscribe(Collections.singletonList("test"));

        while ( true ) {
            // TODO 每隔100毫秒,抓取一次数据--消费者主动拉取数据
            ConsumerRecords<String, String> records =
                    consumer.poll(Duration.ofMillis(100));
            // TODO 打印抓取的数据
            for (ConsumerRecord<String, String> record : records) {
                System.out.println("K = " + record.key() + ", V = " + record.value());
            }
        }
    }
}

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

声明式开发–注解–接收端@KafkaListener注解:指定监听多个主题,指定消费者组id

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

声明式开发-发送端,直接@Resource注入核心类:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

操作分区,偏移量-----偏移量用于定位消息在各个分区内的索引位=====定位一个消息-=》broker号+主题+分区+偏移量

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

默认–拉取消息从最新的偏移量的下一个位置开始读,不指定读取模式是拿不到消息的----特殊—消费者组读取一次后,偏移量会被记住(消费者组id和偏移量时绑定的,读取了两个消息,则偏移量就会移动2),下一次读取从下一位(第3位–空)开始读,即使配置了auto.offset.reset=earliest也没用

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

在配置文件中更改消费者配置:earliest–读取最早的消息:----- 若之前读取过一次,则需要初始化偏移量

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

使用指令更改----偏移量移动到最早的消息上-下一次读取从最早的消息开始读;;;从最新的开始

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

在spring配置文件中更改消费者端读取设置-----偏移量策略:最早的消息开始读

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

生产者发送消息的几种发送消息的对象

messagebuilder组装消息:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

ProducterRecord对象组装消息:Headers可以放一些额外消息;构造方法传参是指定主题和分区和时间戳

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

sendDefault方法发送-----发送时不必再指定发送到哪个主题,而是再配置文件中统一指定发送到的主题,方便管理:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

发送后,返回值对象CompletableFuture–生产者发送消息后,kafka返回一个异步对象CompletableFuture,此时生产者的进程可以继续进行,不阻塞线程进行,(正常情况下要等待kafka发回相应信息后才会继续执行)发送完成后,提供一些方法来获取kafka的返回信息

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

CompletableFuture.get()----获取结果集,阻塞式的,等待结果返回才会继续进行代码

结果集.getXXX()----kafka相应的具体信息–就是刚刚发送的消息的详细

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

非阻塞式:CompletableFuture.thenAccpt()-----非阻塞式获取结果集:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

发送对象类型的消息(常用了)

注入符合泛型的Template对象:第一个键key发送string类型,值value发送字符串“string类型”

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

发送对象类型的消息:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

消息是键值对,键值对都需要序列化,才能通过网络传输到kafka服务器,默认的序列化器只能转化string类型,再spring配置文件中配置生产者–对象类型的数据的序列化器----转化为字节数据byte:====key一般都是字符串,默认的就可以不用改

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

副本机制–Replica-----主副本,从副本—必须在集群部署的环境下才能配置多个副本

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

方法1-命令行方式创建副本:新建一个主题,partition指定分区数量,replication指定副本个数,副本个数不能大于分区数,副本个数不能大于broker节点个数

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

方法2-代码中创建主题的同时指定副本个数-----在一个配置类中@Configuration写–5个分区,一个副本,项目运行自动创建

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

如果多次new topic针对同一个主题,更改部分参数–分区数增加为9,能生效且已经有的消息不会丢失;只能增加不能减小

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

生产者的发送消息的分区策略—不指定分区号,采用一定策略将消息,负载均衡,发送到各个分区中:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

默认策略–根据消息的key值,hash计算出,一个分区号;如果消息没有key值,则随机数%分区数,计算,发送给不同分区

轮询策略–如何配置?需要在自定义配置类中,写一个自定义生产者工厂,或者自己重写工厂方法创建一个自定义的template模板:(或者用上面的自定义一个Producter类也行)

一个map存放自定义生产者的若干配置----PARTITIONER_CLASS_CONFIG常量配置—分区策略–轮询

第一行—kafka服务地址;第二行----键值的序列化

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

生产者工厂加入刚刚的配置map,并根据配置组装一个Template模板,用这个自定义的生产类模板发送消息

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

自定义策略类—了解即可------实现接口

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

实现一个partition方法计算出要发送的目标分区号:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

生产者发送消息流程:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

自定义消息发送前的拦截器—规范消息的键值类型:

实现接口

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

实现1----onsend方法:拦截消息------record对象----消息对象本身----包括键值对,主题,分区,携带头(额外信息),时间戳:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

实现2—该方法用于接收kafka服务器收到消息后的返回信息—正常则返回消息详情metadate,包括偏移量—异常则返回异常信息

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

测试—在创建生产者的配置map中添加–自定义拦截器类:—只需要指定类名,不需要完整类

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

拦截器添加的消息:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

消费者消费数据:被动监听目标主题:

@Payload—只获取消息的消息体

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

@Header–只获取消息的消息头,<可指定获取–请求头携带额外信息–如目标主题,目标分区等等>

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

CosumerRecord类接收完整消息

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

接收对象类型的消息------特殊安全机制,POJO类所在的包不被信任,无法序列化:

监听user对象的消息:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

解决方法:

  1. 显式配置信任包
    在消费者工厂的配置中,添加反序列化器的信任包:

    // Java配置示例
    @Bean
    public ConsumerFactory<String, User> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(JsonDeserializer.TRUSTED_PACKAGES, "com.bjpowernode.model"); // 添加信任包
        // 其他配置...
        return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new JsonDeserializer<>(User.class));
    }
    
  2. 通过配置文件设置
    application.properties中直接指定信任包:

    spring.kafka.consumer.properties.spring.json.trusted.packages=com.bjpowernode.model
    
  3. 动态添加信任包
    如果使用JsonDeserializer反序列化器,可直接调用并指定可信赖的包:

    JsonDeserializer<User> deserializer = new JsonDeserializer<>(User.class);
    deserializer.addTrustedPackages("com.bjpowernode.model");
    

解决方法–手动将user对象转化为json字符后再发送:

工具类—1将任意类型的数据转化为json字符串;2和接收字符串,转化为任意指定类型(泛型)的对象

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

发送时,将对象转化为字符串然后发送—对象类型的字符串数据

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

接收时,调用方法,再将json字符串专户为user对象:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

消息监听器配置----监听到消息后手动向kafka服务器发送一个确认:该消息以消费,需要将偏移量+1.。。。(如果不写则会自动发送确认消息)

配置文件中1:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

被动监听时:多一个手动确认的参数,调用ack方法告知服务器该消息以消费,需要将偏移量加一:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

如果声明参数后,但不调用acknlsge方法,则kafka服务不会添加偏移量—这条信息会被重复消费(偏移量不变)-----场景----监听一条消息后执行业务,当一个业务处理成功无异常时,向kafka发送确认(调用方法),如果异常,则不确认,保证下次还是当前的这条消息:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

指定监听目标-主题,分区,从指定偏移量来监听数据-----全都写在@Listener注解中:

3分区—只当从偏移量3开始读取–

细节理解—默认情况下-即读0,1,2分区的消息读不出,因为当一个消费组消费过一次消息后,分区的偏移量会指向末尾,因此读不到–要么清空偏移量,要么消费组换个名字------而@PartitionOffset–指定分区,并从指定偏移量位开始读,能读到

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

批量消费,一次性取出多条消息:

再配置文件中配置-------batch批量读取----每次最多读取20条

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

接收时使用List接收:ConsumerRecord对象—完整的消息:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

效果–共125条,一次读取20条----size=20:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

消费消息前,拦截器拦截–

需自定义拦截器类—实现…接口:

oncosume()接收消息前执行

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

oncommit()收到消息后,提交确认信息之前—确认信息包括接收后的偏移量:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

自定义配置类中------消费者工厂(consumer…)创建消费者对象,工厂的配置Map中指定–拦截器类类名

配置

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

工厂–使用刚刚的配置

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

自定义配置类中—自定义监听器工厂----注入自定义消费者工厂(consumerfactory),从而使自定义拦截器生效

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

监听注解中,指定自定义监听工厂类:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

效果:接受前,接收后----新增—

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

消费者—消费转发----接受处理后再发送到另一个主题:

消费者接收后–@sendTo注解实现消息转发,转发前,添加一些消息的值

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

多个消费者监听同一个主题—消费时的分区策略:以消费者组为单位去消费

默认分区策略–RangeAssignor–均匀分配:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

在配置类中—更改为轮询策略—消费者工厂配置map中:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

*记得,@Listen注解中要生命自定义的监听器工厂才行 *

轮询效果—纵向的轮询:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

其他的几种分区策略:

也就是说,前者。 会保持原有的分区不变。即使有新的。 呃,消费者进入,或者有原来的消费者离开,也只对新的。 有变化的消费组去更改它的分区。 后者是。 在离开前。 先提前准备好。 顺序的分配。 然后等待消费者离开后。 更精准的将分区分配给下一个消费者。

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

kafka的log文件家中各个文件的作用:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

kafka----consumer_offset分区主题,专门用于协调分区的一个主题:消费者接受消息后,会发回一个offset的确认信息–就是指定当前读到了多少偏移量—每一个offset文件夹对应一个偏移位-----consumer_offset_x中存放第x位的配置信息

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

生产者的offset:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

消费者的offset----确认收到返回一个确认消息–这个过程就是提交本消费者的offset,即偏移量(起始就是索引位)

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

注意:生产者,发送消息之后。那个分区中的。 索引位,也就是所谓的偏移量,就会指定到。 刚刚生产者的offset的下一位(空),所以说你消费者第一次去读的时候也是读不到的,因为消费者默认的起始读取的偏移为就是生产者刚刚的offser位。 他就是生产者已经生产过的消息的下一位。本来就是空。

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

接收者可以重置自己开始读取时的offset—参考最开始的案例

集群部署——kafka真正能处理超大量数据的原因:概念–了解即可

.logs目录下.log文件–存储数据(消息)

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

Broker(节点)就是一个容器—其中存放若干topic(主题)—集群部署起始就是同时部署好几个Brocker

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

概念:

主题-topic-------相当于一个数据库:消息的生产者必须将消息数据发送到某一个主题,而消费者必须从某一个主题中获取消息,并且消费者可以同时消费一个或多个主题的数据。Kafka集群中可以存放多个主题的消息数据。

分区:Partition----------逻辑上将数据库划分

Kafka消息传输采用发布、订阅模式,所以消息生产者必须将数据发送到一个主题,假如发送给这个主题的数据非常多,那么主题所在broker节点的负载和吞吐量就会受到极大的考验,甚至有可能因为热点问题引起broker节点故障,导致服务不可用。一个好的方案就是将一个主题从物理上分成几块,然后将不同的数据块均匀地分配到不同的broker节点上,这样就可以缓解单节点的负载问题。这个主题的分块我们称之为:分区partition。默认情况下,topic主题创建时分区数量为1,也就是一块分区,可以指定参数–partitions改变。Kafka的分区解决了单一主题topic线性扩展的问题,也解决了负载均衡的问题。*

偏移量:Offset:每条消息都有的一个id信息:通过偏移量来定位消息

副本类型---------Leader & Follower----一个集群的.log文件会复制为多个副本—副本的.log文件没有读写的权限—只是用来做备份数据-

分区partition编号–集群形况下–将一个主题中的消息拆分为多个部分–各个部分部署到各个Brocker中,实现消息的并发获取

有点抽象:想想:

场景描述假设我们有一个电商平台,需要展示商品的详细信息页面。这个页面需要从多个服务获取数据,包括:1. 商品基本信息(名称、价格、库存等)2. 用户评价3. 商品图片4. 相关推荐商品在传统的架构中,这些数据可能分别存储在不同的数据库或服务中,前端需要依次调用这些服务来获取数据,然后组合在一起展示给用户。这种方法会导致页面加载时间较长,用户体验不佳。

分区部署–同时取出所有需要用的消息–更快

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

groupid–消费者组–一个功能的消费者,想使用某些消息,但是每个broker中只有部分消息,这就意味着需要访问好多broker才能拿到所有需要的消息,重新连接一个broker需要连接的时间—一个功能的消费者(需要好几个broker中的消息)分组,使得这一个消费者复制成多个–每个消费者只负责一个broker中的一个分区(部分数据)最后再汇总这些消息,避免了来回切换broker连接的时间

也就是所谓的--------每个消费者组中的消费者实例将共同消费订阅主题的所有分区,但每个分区只会被组内的一个消费者实例消费

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

持久化存储–如果一个broker宕机–则这个分区中的部分消息无法被消费了–将一个分区内的消息(.log)复制多份,想要实现一个broker宕机后数据还能恢复,但是一个broker宕机后,下面的所有.log也都失效了,因此将自己的备份.log放到别的broker中,这样自己宕机后,别人能帮你回复

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

备份数据(log文件的复制)–称为副本(follower副本)-只存数据没有读写功能–有一种副本(Loader副本)-可以有读写功能

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

多个broker中存在一个管理者节点(controller)—管理汇报所有的broker的健康状态–每个broker都可以晋升为管理者,当管理者也宕机,会自动选举出一个新的管理者(使用zookeeper实现选举的这一功能)

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

整体集群架构:主副本放在不同的Broker中方便恢复,选举出新的主节点后也能用从副本恢复

主题A,两个分区,三个副本

主题B,一个分区,一个副本

主题C,一个分区,一个副本

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

集群部署:

注意一个点-----生成的uuid是集群的id,要部署的三个节点应使用同样的集群id

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

三个节点的server.properties如下:

process.roles=controller,broker
node.id=2
listeners=PLAINTEXT://:9093,CONTROLLER://:9094
controller.listener.names=CONTROLLER  
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT  
controller.quorum.voters=1@localhost:9092,2@localhost:9094,3@localhost:9096
log.dirs=D:\\kafka_4.0\\kafka_5\\logs
num.partitions=3       
process.roles=controller,broker
node.id=1
listeners=PLAINTEXT://:9091,CONTROLLER://:9092
controller.listener.names=CONTROLLER  
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT  
controller.quorum.voters=1@localhost:9092,2@localhost:9094,3@localhost:9096
log.dirs=D:\\kafka_4.0\\kafka_4\\logs
num.partitions=3       
process.roles=controller,broker
node.id=3
listeners=PLAINTEXT://:9095,CONTROLLER://:9096
controller.listener.names=CONTROLLER  
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT  
controller.quorum.voters=1@localhost:9092,2@localhost:9094,3@localhost:9096
log.dirs=D:\\kafka_4.0\\kafka_6\\logs
num.partitions=3       
按照1,2,3的顺序启动后,1号节点可能会报错:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

但是只要检查端口9096:发现有服务就行,这只是日志打印问题:

netstat -an | findstr 9092

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

工具观察:

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

在springbooot环境中测试连接:自定义创建主题–指定分区数和副本个数—若没有集群则报错:

@Configuration
public class KafkaConfiguration {
    @Bean
    public NewTopic topic1() {
        //指定副本个数为3,若不是集群则会报错
        return new NewTopic("test111", 3, (short)3);
    }
}

效果L:外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

解释:Leader表示当前分区的主副本在id=1的Broker中;replicass表示当前分区的所有副本都在哪个Broker中

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

参数解释:重要

listeners=PLAINTEXT://:9095,CONTROLLER://:9096

其中PLAINTEXT的端口是用于客户端或外界连接的端口;CONTROLLER则是集群内若干节点互相访问的端口,也是选举时要配置的端口

controller.quorum.voters=1@localhost:9092,2@localhost:9094,3@localhost:9096

选举队列----默认1号位为集群中的controller节点--------选举队列中的各个节点的端口是CONTROLLER指定的节点内部通讯端口

其他重要概念:l

ISR副本----只保留,在主副本(Leader)数据修改后,能够及时更新的从副本(follower)

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

LEO----一个数----表示下一位偏移量(索引位)

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

HW—一个偏移量offset类型的数据—在此之前的消息已经同步到从节点,这些消息在从节点中才可以被消费-----LEO-写入了9位,HW只同步了5

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传
new NewTopic(“test111”, 3, (short)3);
}
}


效果L:[外链图片转存中...(img-uHMuaF2g-1744305569647)]

### 解释:Leader表示当前分区的主副本在id=1的Broker中;replicass表示当前分区的所有副本都在哪个Broker中

[外链图片转存中...(img-Jjvvrgi1-1744305569647)]



## 参数解释:重要

***listeners=PLAINTEXT://:9095,CONTROLLER://:9096***

*其中PLAINTEXT的端口是用于客户端或外界连接的端口;CONTROLLER则是集群内若干节点互相访问的端口,也是选举时要配置的端口*

***controller.quorum.voters=1@localhost:9092,2@localhost:9094,3@localhost:9096***

*选举队列----默认1号位为集群中的controller节点--------选举队列中的各个节点的端口是*CONTROLLER指定的节点内部通讯端口



## 其他重要概念:l

### ISR副本----只保留,在主副本(Leader)数据修改后,能够及时更新的从副本(follower)

[外链图片转存中...(img-7nHjq2QO-1744305569647)]

[外链图片转存中...(img-vLsgpCRq-1744305569648)]

### LEO----一个数----表示下一位偏移量(索引位)

[外链图片转存中...(img-U7QrzYgD-1744305569648)]

### HW---一个偏移量offset类型的数据---在此之前的消息已经同步到从节点,这些消息在从节点中才可以被消费-----LEO-写入了9位,HW只同步了5

[外链图片转存中...(img-MGZnxDZB-1744305569648)]

[外链图片转存中...(img-C0bFlCjJ-1744305569648)]
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐