Java技术栈总结:kafka篇

一、# 基础知识

1、安装

  • 部署一台ZooKeeper服务器;
  • 安装jdk;
  • 下载kafka安装包;
  • 上传安装包到kafka服务器上:/usr/local/kafka;
  • 解压缩压缩包;
  • 进入到config目录,修改server.properties配置信息:
#broker.id属性在kafka集群中必须要是唯⼀
broker.id=0#kafka部署的机器ip和提供服务的端⼝号
listeners=PLAINTEXT://192.168.65.60:9092#kafka的消息存储⽂件
log.dir=/usr/local/data/kafka-logs#kafka连接zookeeper的地址
zookeeper.connect=192.168.65.60:2181
  • 进入到bin目录,使用命令启动kafka服务器(带配置文件)
  • ./kafka-server-start.sh -daemon ../config/server.properties
  • 检查kafka是否启动成功:
  • 进入到zk内查看是否有kafka节点:
    /brokers/ids/0

    2、基本概念

名称

说明

Broker

消息中间件处理节点,一个kafka节点为一个broker,一个或者多个broker组成一个kafka集群

Topic

消息主题。kafka根据topic对消息进行分类,发布到kafka集群的每条消息都需要指定一个topic

Producer

消息生产者。向broker发送消息的客户端。

Consumer

消息消费者。从broker读取消息的客户端。

3、主题创建

  • 通过kafka命令向zk中创建一个主题
./kafka-topics.sh --create --zookeeper 172.16.253.35:2181 --replicationfactor 1 --partitions 1 --topic test
  • 查看当前zk中所有的主题
./kafka-topics.sh --list --zookeeper 172.16.253.35:2181 test

4、发送消息

把消息发送给broker的某个topic,打开一个kafka发送消息的客户端,然后开始用客户端向kafka服务器发送消息。

./kafka-console-consumer.sh --bootstrap-server 172.16.253.38:9092 --topic test

5、消费消息

打开一个消费消息的客户端,向kafka服务器的某个主题消费消息。

生产者将消息发送给broker,broker会将消息保存到本地的日志文中。/usr/local/kafka/data/kafka-logs/主题-分区/00000000.log;消息的保存是有序的,通过offset偏移量来描述消息的有序性;消费者消费消息时也是通过offset来描述所要消费消息的位置。

  • 方式一:从当前主题中的最后一条消息的offset + 1 开始消费:
./kafka-console-consumer.sh --bootstrap-server 172.16.253.38:9092 --topic test
  • 方式二:从当前主题的第一条消息开始消费
./kafka-console-consumer.sh --bootstrap-server 172.16.253.38:9092 --from-beginning --topic test

6、单播&&多播消息

如果多个消费者在同一个消费组,那么只有一个消费者可以订阅到topic中的消息。即,同一个消费组中只能有一个消费者收到一个topic中的消息。

不同的消费组订阅同一个topic,那么不同消费组中各只有一个消费者能收到消息。

./kafka-console-consumer.sh --bootstrap-server 172.16.253.38:9092 --consumer-property group.id=testGroup1 --topic test
./kafka-console-consumer.sh --bootstrap-server 172.16.253.38:9092 --consumer-property group.id=testGroup2 --topic test

7、查看消费组信息

/kafka-consumer-groups.sh --bootstrap-server 172.16.253.38:9092 --describe --group testGroup
  • current-offset: 最后被消费的消息的偏移量;
  • Log-end-offset: 消息总量(最后⼀条消息的偏移量);
  • Lag:积压了多少条消息。


二、主题与分区

1、主题 topic

kafka通过topic对消息进行分类,不同的topic会被订阅该topic的消费者消费。

如果一个topic的消息非常多,消息保存在log日志文件中,会占用大量的磁盘空间。为了解决文件过大的问题,kafka提出了Partition分区的概念。

2、分区

通过partition将一个topic中的消息分区来存储。好处:

  • 分区存储,解决了统一存储文件过大的问题;
  • 提升了读写的吞吐量:读和写可以同时在多个分区中进行。

创建多个分区的主题:

./kafka-topics.sh --create --zookeeper 172.16.253.35:2181 --replicationfactor 1 --partitions 2 --topic test1

3、消息日志

  • 00000.log:保存的为消息;
  • __consumer_offset-49:
    • kafka内部创建了主题 “__consumer_offsets” 包含50个分区。这个主题用来存放消费者消费某个主题的偏移量。每个消费者会自己维护消费的主题的偏移量,即每个消费者会把消费的主题的偏移量自主上报给kafka中的默认主题consumer_offsets。kafka为了提升这个主题的并发性,默认设置了50个分区。
      • 提交到哪个分区,通过hash函数确定:
        • hash(cunsumerGroupId)% __consumer_offsets 主题的分区数;
        • 提交到该主题的内容:key为 consumerGroupId+topic+分区号,value为当前的offset值。
  • 文件内容默认保存7天,到期后消息自动删除。

三、集群

kafka的服务端由被称为Broker的服务进程构成,即一个kafka集群由多个Broker组成。如果集群中的某一台机器宕机,其他机器上的Broker仍然能够对外提供服务,确保kafka的高可用性。

1、集群搭建

  • 创建多个server.properties文件
# 0 1 2
broker.id=2
// 9092 9093 9094
listeners=PLAINTEXT://192.168.65.60:9094
// kafka-logs kafka-logs-1 kafka-logs-2
log.dir=/usr/local/data/kafka-logs-2
  • 通过命令分别启动各个broker
./kafka-server-start.sh -daemon ../config/server.properties
./kafka-server-start.sh -daemon ../config/server1.properties
./kafka-server-start.sh -daemon ../config/server2.properties
  • 检查是否启动成功

进入到zk中查看 /brokers/ids 中是否有对应的znode(0,1,2)。

2、副本

# 副本
./kafka-topics.sh --create --zookeeper 172.16.253.35:2181 --replicationfactor 3 --partitions 2 --topic my-replicated-topic
# 查看topic情况
./kafka-topics.sh --describe --zookeeper 172.16.253.35:2181 --topic myreplicated-topic

副本是为了给主题中的分区创建多个备份,多个副本在kafka的集群的多个broker中,会有一个副本作为Leader,其他为Follower。

  • Leader:
    • kafka的写和读操作,都发生在Leader。Leader负责把数据同步给Follower,如果Leader挂了,通过主从选举,从多个Follower中选举产生一个新的Leader。
  • Follower:
    • 接收Leader的数据同步。
  • isr:
    • 可以同步和已经同步的节点会被存入到isr集合中。如果isr中的节点性能较差,会被从isr集合中剔除。

总结:集群中有多个broker,创建主题是可以指明主题有多个分区,可以为分区创建多个副本,不同的副本存放在不同的broker里。

3、集群消费

  • 一个partition(分区)只能被一个消费组中的一个消费者消费,目的是为了保证消费的顺序性。多个partition的多个消费者的消费顺序的顺序性无法得到保证。
  • partition的数量决定了消费组中消费者的数量,同一个消费组中的消费者的数量最好不要超过partition的数量,否则超出的消费者消费不到消息。
  • 如果消费者挂了,会出发rebalance机制,会让其他消费者来消费该分区。

4、Controller、rebalance、Hw

(1)Controller

【Controller选举】:

启动时,每个broker会向zk创建一个临时序号节点,获得序号最小的那个broker将会作为集群中的Controller,负责:

  • Leader选举:当集群中一个副本的Leader挂掉,需要在集群中选举出一个新的Leader,选举从isr集合中的最左边获得。
  • broker信息同步:当集群中有broker新增或者减少,Controller会同步信息给其他broker。
  • 分区信息同步:当集群中有分区新增或者减少,Controller会同步信息给其他broker。

(2)reblance机制

前提:消费组中的消费者没有指明分区来消费;

触发的条件:消费组中的消费者和分区的关系发生变化;

分区分配的策略:reblance之前,分区有三种分配策略:

  • range:根据公式计算每个消费者消费哪几个分区,分区总数/消费者数量 + 1 (根据余数情况确定,前面几个消费者需要“+1”,后面几个不需要)。
  • 轮询:即依次轮着来。
  • sticky:粘合策略。如果需要reblance,会在之前已经分配的基础上进行调整,不会改变之前分配的情况。如果该策略没有开,那么久需要进行整体的重新分配。

(3)HW和LEO

HW是已经完成同步的位置。

消息在写入broker,且每个broker已经完成该消息的同步后,hw才会发生变化。在此之前消费者是消费不到这条消息的。在完成同步后,HW更新后,消费者才能消费到这条消息,这样的目的是为了防止消息丢失

LEO(log-end-offset)是某个副本最后的消息位置。


四、消息的同步异步发送

1、同步发送消息

如果生产者发送消息没有收到ack,生产者会阻塞,阻塞到3s的时间,如果还没有收到消息,会进行重试。重试的次数为3次。

【生产者的三种ack配置】

  • ack=0,kafka-cluster不需要任何的broker收到消息,就立即返回ack给生产者,最容易丢消息,效率最高。
  • ack=1(默认),多副本之间的Leader已经收到消息,并把消息写入到本地的log中,才返回ack给生产者,性能和安全新较为均衡。
  • ack=-1/all,配置min.insync.replicas=2(默认为1,推荐配置大于等于2),此时就需要Leader和一个Follower同步完成后,才返回给ack给生产者(此时集群中有2个broker已经完成数据的接收)。这种方式最安全,但性能最差。

2、异步发送消息

异步发送,生产者发送完消息后就可以执行之后的业务,broker在收到消息后异步调用生产者提供的callback()回调方法。

3、消息发送的缓冲区

  • kafka生产者默认会创建一个消息缓冲区,用来存放要发送的消息,默认为32m;
    • props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
  • kafka本地线程会去缓冲区中一次拉取16k的数据,发送到broker;
    • props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
  • 如果拉取不到16k的数据,间隔10ms也会将已有的数据发送到broker。
    • props.put(ProducerConfig.LINGER_MS_CONFIG, 10);


五、消费者实现

1、消费者自动&&手动提交Offset

(1)提交的内容

“所属的消费组 + 主题 + 分区 + 消费的偏移量”,提交到集群的__consumer_offsets主题里面。

(2)自动提交

消费者poll消息下来后就自动提交offset。

// 是否自动提交offset,默认就是true
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");// 自动提交offset的间隔时间
props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "1000");

注:自动提交会丢消息。因为消费者在消费前提交offset,可能提交完后还没有完成消费,消费者就挂了。

(3)手动提交

需要把自动提交的配置改成false。

props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");

手动提交分为两种:

  • 手动同步提交:在消费完消息后调用同步提交的方法,当集群返回ack前一直阻塞,返回ack后表示提交成功,执行之后的逻辑。
  • 手动异步提交:在消息消费完后提交,不需要等待集群ack,直接执行之后的逻辑,可以设置一个回调方法,供集群调用。

2、长轮询poll消息

(1)默认情况下,消费者一次会拉取500条消息。

props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500);

(2)可以设置长轮询的时间周期,例如1000ms。

  • 如果⼀次poll到500条,就直接执行for循环。
  • 如果这⼀次没有poll到500条。且时间在1秒内,那么⻓轮询继续poll,要么到500条,要么到1s。
  • 如果多次poll都没达到500条,且1秒时间到了,那么直接执行for循环。
  • 如果两次poll的间隔超过30s,集群会认为该消费者的消费能⼒过弱,该消费者被踢出消费组,触发rebalance机制,rebalance机制会造成性能开销。可以通过设置这个参数,让一次poll的消息条数少⼀点。
//⼀次poll最⼤拉取消息的条数,可以根据消费速度的快慢来设置
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500);//如果两次poll的时间如果超出了30s的时间间隔,kafka会认为其消费能⼒过弱,将其踢出消费组。将分区分配给其他消费者。-rebalance
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 30 * 1000);

3、消费者健康状态检查

消费者每隔1skafka集群发送一次心跳,如果集群发现超过10s没有续约的消费者,会将其踢出消费者,触发消费组的reblance机制,将该分区的交给消费组里的其他消费者进行消费。

//consumer给broker发送⼼跳的间隔时间
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 1000);//kafka如果超过10秒没有收到消费者的⼼跳,则会把消费者踢出消费组,进⾏rebalance,把分区分配给其他消费者。
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10 * 1000);

六、常见问题处理

1、防止消息丢失

(1)生产者发送消息到Broker的过程丢失

方式一:异步发送

  • 设置异步发送,发送失败的情况使用回调记录或者重发;
  • 失败重试,配置重试次数。

方式二:同步发送

  • 使用同步发送消息的方式。

(2)消息在Broker中存储丢失

  • 把ack设置为1或者all(-1),设置同步的分区数 >= 2,让Follower节点参与保存数据的确认。

(3)消费者从Broker接收消息丢失

  • 关闭自动提交偏移量,开启手动提交偏移量;
  • 提交方式,把自动提交改成手动提交(最好使用 同步 + 异步 提交)。

2、防止重复消费

如果生产者发送消息后,由于网络抖动等问题,没有收到ack,但是实际上broker已经收到了消息。此时,生产者会进行重试,于是broker就会收到多条相同的消息,从而造成重复消费。

解决:

  • 生产者关闭重试。这种方式会造成消息丢失(不推荐);
  • 消费者关闭自动提交偏移量,开启手动提交偏移量;
  • 消费者解决非幂等性消费问题:
    • 在数据库中创建联合主键,防止相同的主键创建出多条记录。
    • 使用分布式锁,以业务id为锁。保证只有一条记录能够创建成功。

3、保证顺序性消费

问题原因:一个topic的数据可能存储在不同的分区中,每个分区都有一个按照顺序存储的偏移量。如果消费者关联了多个分区,则不能保证顺序性。

解决该问题,只需要保证需要顺序消费的消息出现在同一个分区。

解决方法:

  • 方式一:
    • 发送消息时,指定分区号;
    • 发送消息时,按照相同的业务设置相同的key(默认情况下,分区是通过key的hashcode值来确定分区的。因此,key一样的话,分区也是一样的);
  • 方式二(不推荐):
    • 生产者:使用同步发送,ack设置成非0的值(1或者-1(all))。
    • 消费者:主题只设置一个分区,消费组只设置一个消费者。

主:实际kafka顺序消费的场景不多,因为会牺牲掉性能。

4、消息积压

(1)出现的原因

消费者的消费速度赶不上生产者的生产速度,导致kafka中大量的数据没有被消费。

随着积压消息的增多,消费者的寻址性能会下降,最终导致整个kafka对外提供服务的性能很差,从而造成其他服务访问速度变慢,造成服务雪崩。

(2)解决方案

  • 消费者中,使用多线程,充分利用机器的性能进行消费消息。
  • 通过业务的架构设计,提升业务层面消费的性能。
  • 创建多个消费组,多个消费者,部署到其他机器上,一起消费,提高消费者的消费速度。
  • 创建一个消费者,该消费者kafka另建一个主题,配上多个分区,多个分区再配上多个消费者。该消费者将消息poll下来,不进行消费,直接转发到新建的主题上。此时,新的主题的多个分区的多个消费者就开始一起消费了。(不常用)


参考:

https://www.bilibili.com/video/BV1Xy4y1G7zA?p=1;

https://www.bilibili.com/video/BV1yT411H7YK

https://www.jianshu.com/p/d3e963ff8b70;

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.mzph.cn/bicheng/41955.shtml

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

Buuctf之SimpleRev做法

首先,查个壳,64bit,那就丢进ida64中进行反编译进来之后,我们进入main函数,发现里面没什么东西,那就shiftf12搜索字符串,找到关键字符串,双击进入然后再选中该字符串,ctrl…

Python爬取股票信息-并进行数据可视化分析,绘股票成交量柱状图

为了使用Python爬取股票信息并进行数据可视化分析,我们可以使用几个流行的库:requests 用于网络请求,pandas 用于数据处理,以及 matplotlib 或 seaborn 用于数据可视化。 步骤 1: 安装必要的库 首先,确保安装了以下P…

【C语言】指针(1):入门理解篇

目录 一、内存和地址 1.1内存 1.2 深入理解计算机编址 二、指针变量和地址 2.1 取地址操作符(&) 2.2 指针变量和解应用操作符 2.2.1 指针变量 2.2.2 解引用操作符 2.3指针变量的大小 三、指针变量类型的意义 3.1 指针的解引用 3.1指针-整数…

Micron近期发布了32Gb DDR5 DRAM

Micron Technology近期发布了一项内存技术的重大突破——一款32Gb DDR5 DRAM芯片,这项创新不仅将存储容量翻倍,还显著提升了针对人工智能(AI)、机器学习(ML)、高性能计算(HPC)以及数…

2024年最新运维面试题(附答案)

作者简介:一名云计算网络运维人员、每天分享网络与运维的技术与干货。 公众号:网络豆云计算学堂 座右铭:低头赶路,敬事如仪 个人主页: 网络豆的主页​​​​​ 一.选择题 1.HTTP协议默认使用哪个端口…

科普文:构建可扩展的微服务架构设计方案

前言 微服务架构是一种新兴的软件架构风格,它将单个应用程序拆分成多个小的服务,每个服务都运行在自己的进程中,这些服务通过网络进行通信。这种架构的优势在于它可以提高应用程序的可扩展性、可维护性和可靠性。 在传统的应用程序架构中&…

高效管理个人日程,智慧校园行政办公全指南

在智慧校园的行政办公体系里,个人日程管理功能担当起协调与优化每位教职员工日常安排的角色,它像一位贴心的时间助理,确保工作与私人生活的和谐并进。这一功能设计得既直观又灵活,让使用者能以自己偏好的视角审视时间规划&#xf…

创新配置,秒级采集,火爆短视频评论抓取

快速采集评论数据的好处 快速采集评论数据是在当今数字信息时代的市场趋势分析和用户反馈分析中至关重要的环节。通过准确获取并分析大量用户评论,您将能够更好地了解消费者的需求、情感和偏好。集蜂云采集平台提供了一种简单配置的方法,使您能够快速采…

Deep Filtered Back Projection for CT Reconstruction

CT重建中的深度滤波反投影 论文链接:https://ieeexplore.ieee.org/document/10411896 项目链接: ABSTRACT 滤波反投影(FBP)是一种经典的计算机断层扫描(CT)重建解析算法,具有很高的计算效率。然而,用FBP重建的图像往往存在过多…

NATAPP内网穿透使用

1. natapp能干嘛 可以将本地的内网ip映射到外网上,远程访问该连接,实现外网展示网站。平时做的应用开发都只能在局域网本地访问,通过内网穿透,可以通过外网进行访问。 2. 注册用户 网址:https://natapp.cn/自行完成…

什么是 Elasticsearch 数据预热?

引言:在现代的信息检索和数据分析领域,Elasticsearch 已经成为一个广泛应用的分布式搜索和分析引擎。作为开源项目的一部分,Elasticsearch 提供了强大的实时搜索和分析能力,使得处理大规模数据变得更加高效和可靠。然而&#xff0…

Canary,三种优雅姿势绕过

Canary(金丝雀),栈溢出保护 canary保护是防止栈溢出的一种措施,其在调用函数时,在栈帧的上方放入一个随机值 ,绕过canary时首先需要泄漏这个随机值,然后再钩爪ROP链时将其作为垃圾数据写入&…

对接海康sdk-linux下复制jar包中resource目录的文件夹

背景 在集成海康sdk时,需要将一些组件放到项目中作为静态资源,并且海康的sdk初始化也需要加载这些静态资源,在windows下,使用一些File路径的方式是可以正确加载的,但是在linux上就会加载失败。 首先我是将海康的sdk组件放到resource下的,并且按照windows和linux设置了两…

轻松快速上手Thekey库,实现数据加密无忧

Thekey的概述: Thekey库是一个Python库,旨在简化数据加密、解密、签名和验证的过程。它提供了一套简洁易用的接口,用于处理各种加密任务,适合需要在应用程序中实现安全数据处理的开发人员. 安装Thekey库 pip install thekey使用Thekey库进行基本加密和解密操作的…

【笔记】TimEP Safety Mechanisms方法论

1.TimEPM Overview 三大监控方法: Alive Supervision 实时监督Logical Supervision 逻辑监督Deadline Supervision 限时监督相关模块框图: 相关模块调用框图: 每个MCU核开启内狗(1核1狗),内狗用于监控相应核的TASK超时,超时后软reset MCU内狗时钟需要独立于OS时钟,两…

C++下Protobuf学习

C下Protobuf简单学习 Protobuf(Protocol Buffers)协议是一种由 Google 开发的高效的、跨语言的、平台无关的数据序列化协议,提供二进制序列化格式和相关的技术,它用于高效地序列化和反序列化结构化数据,通常用于网络通…

DDR3(三)

目录 1 预取1.1 什么是预取1.2 预取有哪些好处1.3 结构框图1.4 总结 2 突发2.1 什么是突发2.2 突发与预取 本文讲解DDR中常见的两个术语:预取和突发,对这两个概念理解的关键在于地址线的低位是否参与译码,具体内容请继续往下看。 1 预取 1.1…

JDBC【封装工具类、SQL注入问题】

day54 JDBC 封装工具类01 创建配置文件 DBConfig.properties driverNamecom.mysql.cj.jdbc.Driver urljdbc:mysql://localhost:3306/qnz01?characterEncodingutf8&serverTimezoneUTC usernameroot passwordroot新建配置文件,不用写后缀名 创建工具类 将变…

C++笔试强训2

文章目录 一、选择题二、编程题 一、选择题 和笔试强训1的知识点考的一样,因为输出的是double类型所以后缀为f,m.n对其30个字符所以m是30,精度是4所以n是4,不加符号默认是右对齐,左对齐的话前面加-号,所以答案是-30.4f…

推荐Bulk Image Downloader插件下载网页中图片链接很好用

推荐:Bulk Image Downloader chome浏览器插件下载图片链接,很好用。 有个网页,上面放了数千的gif的电路图,手工下载会累瘫了不可。想找一个工具分析它的静态链接并下载,找了很多推荐的下载工具,都是不能分…