Kafka的基本使用及環(huán)境安裝
認(rèn)識(shí)Kafka
消息隊(duì)列
消息隊(duì)列是分布式系統(tǒng)和現(xiàn)代應(yīng)用架構(gòu)中至關(guān)重要的中間件。它的核心作用是解耦、異步和削峰填谷,像一個(gè)高效的“通信員”和“緩沖池”協(xié)調(diào)不同組件之間的工作。
消息隊(duì)列的核心概念
生產(chǎn)者: 產(chǎn)生消息(數(shù)據(jù)、任務(wù)請(qǐng)求、事件通知)并發(fā)送到隊(duì)列的應(yīng)用程序或服務(wù)。
消息隊(duì)列: 一個(gè)臨時(shí)的、持久化的存儲(chǔ)區(qū)域(通常基于內(nèi)存、磁盤或數(shù)據(jù)庫(kù)),用于存放生產(chǎn)者發(fā)送的消息。消息按照先進(jìn)先出的順序存儲(chǔ),但很多隊(duì)列支持優(yōu)先級(jí)、延遲等特性。
消費(fèi)者: 從隊(duì)列中獲取消息并進(jìn)行處理的應(yīng)用程序或服務(wù)。
消息: 隊(duì)列中傳輸?shù)臄?shù)據(jù)單元,通常包含有效載荷(實(shí)際數(shù)據(jù))和元數(shù)據(jù)(如ID、時(shí)間戳、優(yōu)先級(jí)等)。
核心價(jià)值與解決的問題
解耦:
問題: 系統(tǒng)組件(服務(wù))之間直接調(diào)用會(huì)導(dǎo)致緊密耦合。一個(gè)組件的變更、故障或性能瓶頸會(huì)直接影響其他依賴它的組件。擴(kuò)展也變得困難。
解決: 生產(chǎn)者只需將消息發(fā)送到隊(duì)列,無需知道誰(消費(fèi)者)會(huì)處理它,消費(fèi)者只需從隊(duì)列訂閱消息,無需知道消息是誰(生產(chǎn)者)發(fā)送的。雙方只依賴隊(duì)列,不直接依賴對(duì)方,大大降低了耦合度。系統(tǒng)更靈活、更易于維護(hù)和擴(kuò)展。
異步:
問題: 同步調(diào)用要求調(diào)用方(生產(chǎn)者)必須等待被調(diào)用方(消費(fèi)者)處理完成并返回結(jié)果才能繼續(xù)執(zhí)行。如果處理耗時(shí)很長(zhǎng),調(diào)用方會(huì)被阻塞,資源利用率低,用戶體驗(yàn)差(如網(wǎng)頁卡頓)。
解決: 生產(chǎn)者發(fā)送消息到隊(duì)列后即可返回,無需等待消費(fèi)者處理。消費(fèi)者在后臺(tái)異步地從隊(duì)列拉取消息進(jìn)行處理。這顯著提高了系統(tǒng)的吞吐量和響應(yīng)速度。
削峰填谷:
問題: 系統(tǒng)流量往往存在高峰和低谷。高峰期如果請(qǐng)求量遠(yuǎn)超消費(fèi)者處理能力,會(huì)導(dǎo)致系統(tǒng)過載、崩潰或請(qǐng)求超時(shí)。低谷期資源又可能閑置。
解決: 隊(duì)列作為緩沖區(qū),在流量高峰時(shí)積壓請(qǐng)求,平滑地將大量請(qǐng)求暫存起來。消費(fèi)者按照自己的穩(wěn)定處理能力從隊(duì)列中拉取消息進(jìn)行處理,避免了瞬間洪峰壓垮下游系統(tǒng)。在流量低谷時(shí),消費(fèi)者可以繼續(xù)處理隊(duì)列中積壓的消息。
冗余與可靠性:
問題: 直接調(diào)用時(shí),如果消費(fèi)者臨時(shí)不可用(故障、重啟、維護(hù)),生產(chǎn)者的請(qǐng)求會(huì)丟失或失敗。
解決: 消息隊(duì)列通常提供消息持久化功能(將消息寫入磁盤)。即使消費(fèi)者暫時(shí)離線,消息也會(huì)安全存儲(chǔ)在隊(duì)列中,待消費(fèi)者恢復(fù)后繼續(xù)處理,確保消息不丟失。許多隊(duì)列還提供確認(rèn)機(jī)制(ACK),消費(fèi)者處理成功后才會(huì)從隊(duì)列中移除消息。
可伸縮性:
問題: 單一消費(fèi)者處理能力有限,難以應(yīng)對(duì)增長(zhǎng)的業(yè)務(wù)量。
解決: 可以很容易地增加消費(fèi)者的數(shù)量(水平擴(kuò)展),讓多個(gè)消費(fèi)者并行地從同一個(gè)隊(duì)列中拉取消息進(jìn)行處理,顯著提高系統(tǒng)的整體吞吐量。隊(duì)列本身也可以做成分布式集群來應(yīng)對(duì)高吞吐量需求。
順序保證:
問題: 在分布式環(huán)境中保證消息處理的嚴(yán)格順序很困難。
解決: 雖然完全全局有序很難,但許多消息隊(duì)列能保證分區(qū)有序或隊(duì)列有序(在單個(gè)隊(duì)列/分區(qū)內(nèi),消息按照發(fā)送順序被消費(fèi))。這對(duì)于某些需要保證因果關(guān)系的業(yè)務(wù)場(chǎng)景(如賬戶流水)非常重要。
緩沖:
問題: 生產(chǎn)者和消費(fèi)者的處理速度不一致。
解決: 隊(duì)列天然提供了緩沖能力,允許生產(chǎn)者和消費(fèi)者以各自不同的速率工作,不會(huì)互相拖累。
常見的消息隊(duì)列有RabbitMQ,Kafka,RocketMQ。這里主要介紹Kafka。
Kafka
Kafka 通常指 Apache Kafka,這是一個(gè)開源的、分布式的、高吞吐量、低延遲的流處理平臺(tái)。它最初由 LinkedIn 開發(fā),后來捐贈(zèng)給了 Apache 軟件基金會(huì),并迅速成為大數(shù)據(jù)和實(shí)時(shí)數(shù)據(jù)處理領(lǐng)域的核心基礎(chǔ)設(shè)施之一。
Kafka 不僅僅是一個(gè)消息隊(duì)列,它是一個(gè)高吞吐、低延遲、分布式、持久化、可水平擴(kuò)展的流數(shù)據(jù)平臺(tái)。它設(shè)計(jì)之初就是為了處理持續(xù)產(chǎn)生、體量巨大、需要實(shí)時(shí)處理的“數(shù)據(jù)流”。
ZooKeeper是一個(gè)開源的分布式應(yīng)用程序協(xié)調(diào)軟件,而Kafka是分布式事件處理平臺(tái),底層是使用分布式架構(gòu)設(shè)計(jì),所以Kafka的多個(gè)節(jié)點(diǎn)之間是采用zookeeper來實(shí)現(xiàn)協(xié)調(diào)調(diào)度的。
ZooKeeper
ZooKeeper是一個(gè)開源的分布式應(yīng)用程序協(xié)調(diào)軟件,而Kafka是分布式事件處理平臺(tái),底層是使用分布式架構(gòu)設(shè)計(jì),所以Kafka的多個(gè)節(jié)點(diǎn)之間是采用zookeeper來實(shí)現(xiàn)協(xié)調(diào)調(diào)度的。
Zookeeper的核心作用
ZooKeeper的數(shù)據(jù)存儲(chǔ)結(jié)構(gòu)可以簡(jiǎn)單地理解為一個(gè)Tree結(jié)構(gòu),而Tree結(jié)構(gòu)上的每一個(gè)節(jié)點(diǎn)可以用于存儲(chǔ)數(shù)據(jù),所以一般情況下,我們可以將分布式系統(tǒng)的元數(shù)據(jù)(環(huán)境信息以及系統(tǒng)配置信息)保存在ZooKeeper節(jié)點(diǎn)中。
ZooKeeper創(chuàng)建數(shù)據(jù)節(jié)點(diǎn)時(shí),會(huì)根據(jù)業(yè)務(wù)場(chǎng)景創(chuàng)建臨時(shí)節(jié)點(diǎn)或永久(持久)節(jié)點(diǎn)。永久節(jié)點(diǎn)就是無論客戶端是否連接上ZooKeeper都一直存在的節(jié)點(diǎn),而臨時(shí)節(jié)點(diǎn)指的是客戶端連接時(shí)創(chuàng)建,斷開連接后刪除的節(jié)點(diǎn)。同時(shí),ZooKeeper也提供了Watch(監(jiān)控)機(jī)制用于監(jiān)控節(jié)點(diǎn)的變化,然后通知對(duì)應(yīng)的客戶端進(jìn)行相應(yīng)的變化。Kafka軟件中就內(nèi)置了ZooKeeper的客戶端,用于進(jìn)行ZooKeeper的連接和通信。
Kafka的基本使用
環(huán)境安裝
我們這里先安裝簡(jiǎn)單的Windows單機(jī)環(huán)境。在安裝之前務(wù)必先安裝Java8。
下載Kafka:Kafka下載地址Apache Kafka: A Distributed Streaming Platform.
https://kafka.apache.org/downloads
選擇版本為2.13-3.8.0

下載完成后進(jìn)行解壓,解壓目錄放在非系統(tǒng)盤根目錄下。為了訪問方便,可以將解壓后的文件夾名稱修改為Kafka
Kafka的文件目錄
bin | linux系統(tǒng)下可執(zhí)行腳本文件 |
bin/windows | windows系統(tǒng)下可執(zhí)行腳本文件 |
config | 配置文件 |
libs | 依賴類庫(kù) |
licenses | 許可信息 |
site-docs | 文檔 |
logs | 服務(wù)日志 |
啟動(dòng)zookeeper
當(dāng)前版本的Kafka軟件仍然依賴Zookeeper,所以啟動(dòng)Kafka之前,需要先啟動(dòng)Zookeeper,Kafka軟件內(nèi)置了Zookeeper,所以無需額外安裝,直接調(diào)用啟動(dòng)腳本即可。
1. 進(jìn)入Kafka解壓縮文件夾的config目錄,修改zookeeper.properties配置文件
修改dataDir配置,用于設(shè)置ZooKeeper數(shù)據(jù)存儲(chǔ)位置,該路徑如果不存在會(huì)自動(dòng)創(chuàng)建。
dataDir=D:/kafka/data/zk

在kafka解壓縮后的目錄中創(chuàng)建Zookeeper啟動(dòng)腳本文件:zk.cmd。
輸入:
call bin/windows/zookeeper-server-start.bat config/zookeeper.properties
上述指令就是調(diào)用zookeeper啟動(dòng)命令,同時(shí)指定配置文件
雙擊啟動(dòng)即可:

啟動(dòng)完成。
啟動(dòng)Kafka
進(jìn)入Kafka解壓縮文件夾的config目錄,修改server.properties配置文件.

設(shè)置Kafka數(shù)據(jù)的存儲(chǔ)目錄。如果文件目錄不存在,會(huì)自動(dòng)生成。
在kafka解壓縮后的目錄中創(chuàng)建Kafka啟動(dòng)腳本文件:kfk.cmd。
輸入:
call bin/windows/kafka-server-start.bat config/server.properties
雙擊啟動(dòng)即可:

DOS窗口中,輸入jps指令,查看當(dāng)前啟動(dòng)的軟件進(jìn)程:

這里名稱為QuorumPeerMain的就是ZooKeeper軟件進(jìn)程,名稱為Kafka的就是Kafka系統(tǒng)進(jìn)程。此時(shí),說明Kafka已經(jīng)可以正常使用了。
消息主題
在發(fā)布訂閱模型中,為了讓消費(fèi)者對(duì)感興趣的消息進(jìn)行消費(fèi),而不是消費(fèi)所有消息,所以就定義了主題(Topic),也就是說將不同的消息進(jìn)行分類,分成不同的主題(Topic),然后消息生產(chǎn)者在生成消息時(shí),就會(huì)向指定的主題(Topic)中發(fā)送,而消息消費(fèi)者也可以訂閱自己感興趣的主題(Topic)并從中獲取消息。
有很多種方式都可以操作Kafka消息中的主題(Topic):命令行、第三方工具、Java API、自動(dòng)創(chuàng)建。而對(duì)于初學(xué)者來講,掌握基本的命令行操作是必要的。所以接下來,我們采用命令行進(jìn)行操作。
創(chuàng)建主題
使用命令行方式創(chuàng)建主題test
打開DOS窗口,在確保Zookeeper和Kafkfa啟動(dòng)的情況下,進(jìn)入Kafkfa解壓目錄下的bin/windows目錄。
輸入如下命令創(chuàng)建主題test: kafka-topics.bat --bootstrap-server localhost:9092 --create --topic test

test主題創(chuàng)建完成。
查詢主題
輸入如下命令進(jìn)行主題查詢:kafka-topics.bat --bootstrap-server localhost:9092 --list

修改主題
kafka-topics.bat --bootstrap-server localhost:9092 --topic test --alter --partitions 2
上述命令將test主題的分區(qū)數(shù)量設(shè)置為2.關(guān)于分區(qū)的信息,后面會(huì)詳細(xì)介紹。
發(fā)送數(shù)據(jù)
命令行操作
使用命令行方式發(fā)送:
kafka-console-producer.bat --bootstrap-server localhost:9092 --topic test

上述操作就是在控制臺(tái)生成數(shù)據(jù),hello kafka 這里的數(shù)據(jù)需要回車,才會(huì)發(fā)送到Kafka服務(wù)器。
JavaAPI操作
引入依賴
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.8.0</version>
</dependency>編寫生產(chǎn)者
public class ProducerTest {
public static void main(String[] args) {
// 配置屬性集合
Map<String, Object> configMap = new HashMap<>();
// 配置屬性:Kafka服務(wù)器集群地址
configMap.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
// 配置屬性:Kafka生產(chǎn)的數(shù)據(jù)為KV對(duì),所以在生產(chǎn)數(shù)據(jù)進(jìn)行傳輸前需要分別對(duì)K,V進(jìn)行對(duì)應(yīng)的序列化操作
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");
// 創(chuàng)建Kafka生產(chǎn)者對(duì)象,建立Kafka連接
// 構(gòu)造對(duì)象時(shí),需要傳遞配置參數(shù)
KafkaProducer<String, String> producer = new KafkaProducer<>(configMap);
// 準(zhǔn)備數(shù)據(jù),定義泛型
// 構(gòu)造對(duì)象時(shí)需要傳遞 【Topic主題名稱】,【Key】,【Value】三個(gè)參數(shù)
for (int i = 0; i < 10; i++) {
ProducerRecord<String, String> record = new ProducerRecord<String, String>(
"test", "key" + i, "value" + i
);
// 生產(chǎn)(發(fā)送)數(shù)據(jù)
producer.send(record);
}
// 關(guān)閉生產(chǎn)者連接
producer.close();
}
}消費(fèi)數(shù)據(jù)
命令行操作
kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic test --from-beginning

JavaAPI操作
public class ConsumerTest {
public static void main(String[] args) {
// 創(chuàng)建配置對(duì)象
Map<String, Object> configMap = new HashMap<>();
configMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
// 反序列化類配置
configMap.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
configMap.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
// 組ID配置
configMap.put(ConsumerConfig.GROUP_ID_CONFIG, "test");
// 創(chuàng)建消費(fèi)者對(duì)象
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(configMap);
// 從kafka主題中獲取對(duì)象 訂閱主題
consumer.subscribe(Collections.singleton("test"));
// 消費(fèi)者從Kafka主題中拉取數(shù)據(jù)
while (true) {
ConsumerRecords<String, String> datas = consumer.poll(100);
for (ConsumerRecord<String, String> data : datas) {
System.out.println(data);
}
}
// 關(guān)閉消費(fèi)者對(duì)象
// consumer.close();
}
}到此這篇關(guān)于Kafka的基本使用的文章就介紹到這了,更多相關(guān)Kafka使用內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- springboot使用kafka的過程
- Python使用Apache Kafka時(shí)Poll拉取速度慢的解決方法
- SpringBoot使用Kafka來優(yōu)化接口請(qǐng)求的并發(fā)方式
- 如何使用Apache Kafka 構(gòu)建實(shí)時(shí)數(shù)據(jù)處理應(yīng)用
- springboot使用kafka事務(wù)的示例代碼
- Spring Kafka中@KafkaListener注解的參數(shù)與使用小結(jié)
- springboot使用@KafkaListener監(jiān)聽多個(gè)kafka配置實(shí)現(xiàn)
- springboot連接kafka集群的使用示例
- spring?kafka?@KafkaListener詳解與使用過程
相關(guān)文章
MyBatis注解方式之@Update/@Delete使用詳解
這篇文章主要介紹了MyBatis注解方式之@Update/@Delete使用詳解,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過來看看吧2020-11-11
Mybatis-Plus雪花id的使用以及解析機(jī)器ID和數(shù)據(jù)標(biāo)識(shí)ID實(shí)現(xiàn)
這篇文章主要介紹了Mybatis-Plus雪花id的使用以及解析機(jī)器ID和數(shù)據(jù)標(biāo)識(shí)ID實(shí)現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2020-08-08
一文帶你掌握J(rèn)ava?LinkedBlockingQueue
LinkedBlockingQueue?是一個(gè)可選有界阻塞隊(duì)列,這篇文章主要為大家詳細(xì)介紹了Java中LinkedBlockingQueue的實(shí)現(xiàn)原理與適用場(chǎng)景,感興趣的可以了解一下2023-04-04
減小Maven項(xiàng)目生成的JAR包體積實(shí)現(xiàn)提升運(yùn)維效率
在Maven構(gòu)建Java項(xiàng)目過程中,減小JAR包體積可通過排除不必要的依賴和使依賴jar包獨(dú)立于應(yīng)用jar包來實(shí)現(xiàn),在pom.xml文件中使用<exclusions>標(biāo)簽排除不需要的依賴,有助于顯著降低JAR包大小,此外,將依賴打包到應(yīng)用外,可減少應(yīng)用包的體積2024-10-10
springboot中JetCache的使用方法小結(jié)
本文主要介紹了springboot中JetCache的使用方法小結(jié),文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2025-10-10
SpringBoot + Shiro前后端分離權(quán)限
這篇文章主要為大家詳細(xì)介紹了SpringBoot + Shiro前后端分離權(quán)限,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2019-12-12

