Java分布式學(xué)習(xí)之Kafka消息隊(duì)列
介紹
Apache Kafka 是分布式發(fā)布-訂閱消息系統(tǒng),在 kafka官網(wǎng)上對(duì) kafka 的定義:一個(gè)分布式發(fā)布-訂閱消息傳遞系統(tǒng)。 它最初由LinkedIn公司開發(fā),Linkedin于2010年貢獻(xiàn)給了Apache基金會(huì)并成為頂級(jí)開源項(xiàng)目。Kafka是一種快速、可擴(kuò)展的、設(shè)計(jì)內(nèi)在就是分布式的,分區(qū)的和可復(fù)制的提交日志服務(wù)。
注意:Kafka并沒(méi)有遵循JMS規(guī)范(),它只提供了發(fā)布和訂閱通訊方式。
kafka中文官網(wǎng):http://kafka.apachecn.org/quickstart.html
Kafka核心相關(guān)名稱
- Broker:Kafka節(jié)點(diǎn),一個(gè)Kafka節(jié)點(diǎn)就是一個(gè)broker,多個(gè)broker可以組成一個(gè)Kafka集群
- Topic:一類消息,消息存放的目錄即主題,例如page view日志、click日志等都可以以topic的形式存在,Kafka集群能夠同時(shí)負(fù)責(zé)多個(gè)topic的分發(fā)
- massage: Kafka中最基本的傳遞對(duì)象。
- Partition:topic物理上的分組,一個(gè)topic可以分為多個(gè)partition,每個(gè)partition是一個(gè)有序的隊(duì)列。Kafka里面實(shí)現(xiàn)分區(qū),一個(gè)broker就是表示一個(gè)區(qū)域。
- Segment:partition物理上由多個(gè)segment組成,每個(gè)Segment存著message信息
- Producer : 生產(chǎn)者,生產(chǎn)message發(fā)送到topic
- Consumer : 消費(fèi)者,訂閱topic并消費(fèi)message, consumer作為一個(gè)線程來(lái)消費(fèi)
- Consumer Group:消費(fèi)者組,一個(gè)Consumer Group包含多個(gè)consumer
- Offset:偏移量,理解為消息 partition 中消息的索引位置
主題和隊(duì)列的區(qū)別:
隊(duì)列是一個(gè)數(shù)據(jù)結(jié)構(gòu),遵循先進(jìn)先出原則
kafka集群安裝
參考官方文檔:https://kafka.apachecn.org/quickstart.html
- 每臺(tái)服務(wù)器上安裝jdk1.8環(huán)境
- 安裝Zookeeper集群環(huán)境
- 安裝kafka集群環(huán)境
- 運(yùn)行環(huán)境測(cè)試

安裝jdk環(huán)境和zookeeper這里不詳述了。
kafka為什么依賴于zookeeper:kafka會(huì)將mq信息存放到zookeeper上,為了使整個(gè)集群能夠方便擴(kuò)展,采用zookeeper的事件通知相互感知。
kafka集群安裝步驟:
1、下載kafka的壓縮包,下載地址:https://kafka.apachecn.org/downloads.html
2、解壓安裝包
tar -zxvf kafka_2.11-1.0.0.tgz
3、修改kafka的配置文件 config/server.properties
配置文件修改內(nèi)容:
- zookeeper連接地址:
zookeeper.connect=192.168.1.19:2181 - 監(jiān)聽的ip,修改為本機(jī)的ip
listeners=PLAINTEXT://192.168.1.19:9092 - kafka的brokerid,每臺(tái)broker的id都不一樣
broker.id=0
4、依次啟動(dòng)kafka
./kafka-server-start.sh -daemon config/server.properties
kafka使用
kafka文件存儲(chǔ)
topic是邏輯上的概念,而partition是物理上的概念,每個(gè)partition對(duì)應(yīng)于一個(gè)log文件,該log文件中存儲(chǔ)的就是Producer生成的數(shù)據(jù)。Producer生成的數(shù)據(jù)會(huì)被不斷追加到該log文件末端,為防止log文件過(guò)大導(dǎo)致數(shù)據(jù)定位效率低下,Kafka采取了分片和索引機(jī)制,將每個(gè)partition分為多個(gè)segment,每個(gè)segment包括:“.index”文件、“.log”文件和.timeindex等文件。這些文件位于一個(gè)文件夾下,該文件夾的命名規(guī)則為:topic名稱+分區(qū)序號(hào)。
例如:執(zhí)行命令新建一個(gè)主題,分三個(gè)區(qū)存放放在三個(gè)broker中:
./kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic kaico


- 一個(gè)partition分為多個(gè)segment
- .log 日志文件
- .index 偏移量索引文件
- .timeindex 時(shí)間戳索引文件
- 其他文件(partition.metadata,leader-epoch-checkpoint)
Springboot整合kafka
maven依賴
<dependencies>
<!-- springBoot集成kafka -->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<!-- SpringBoot整合Web組件 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>
yml配置
# kafka
spring:
kafka:
# kafka服務(wù)器地址(可以多個(gè))
# bootstrap-servers: 192.168.212.164:9092,192.168.212.167:9092,192.168.212.168:9092
bootstrap-servers: www.kaicostudy.com:9092,www.kaicostudy.com:9093,www.kaicostudy.com:9094
consumer:
# 指定一個(gè)默認(rèn)的組名
group-id: kafkaGroup1
# earliest:當(dāng)各分區(qū)下有已提交的offset時(shí),從提交的offset開始消費(fèi);無(wú)提交的offset時(shí),從頭開始消費(fèi)
# latest:當(dāng)各分區(qū)下有已提交的offset時(shí),從提交的offset開始消費(fèi);無(wú)提交的offset時(shí),消費(fèi)新產(chǎn)生的該分區(qū)下的數(shù)據(jù)
# none:topic各分區(qū)都存在已提交的offset時(shí),從offset后開始消費(fèi);只要有一個(gè)分區(qū)不存在已提交的offset,則拋出異常
auto-offset-reset: earliest
# key/value的反序列化
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
producer:
# key/value的序列化
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
# 批量抓取
batch-size: 65536
# 緩存容量
buffer-memory: 524288
# 服務(wù)器地址
bootstrap-servers: www.kaicostudy.com:9092,www.kaicostudy.com:9093,www.kaicostudy.com:9094
生產(chǎn)者
@RestController
public class KafkaController {
/**
* 注入kafkaTemplate
*/
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
/**
* 發(fā)送消息的方法
*
* @param key
* 推送數(shù)據(jù)的key
* @param data
* 推送數(shù)據(jù)的data
*/
private void send(String key, String data) {
// topic 名稱 key data 消息數(shù)據(jù)
kafkaTemplate.send("kaico", key, data);
}
// test 主題 1 my_test 3
@RequestMapping("/kafka")
public String testKafka() {
int iMax = 6;
for (int i = 1; i < iMax; i++) {
send("key" + i, "data" + i);
}
return "success";
}
}消費(fèi)者
@Component
public class TopicKaicoConsumer {
/**
* 消費(fèi)者使用日志打印消息
*/
@KafkaListener(topics = "kaico") //監(jiān)聽的主題
public void receive(ConsumerRecord<?, ?> consumer) {
System.out.println("topic名稱:" + consumer.topic() + ",key:" +
consumer.key() + "," +
"分區(qū)位置:" + consumer.partition()
+ ", 下標(biāo)" + consumer.offset());
//輸出key對(duì)應(yīng)的value的值
System.out.println(consumer.value());
}
}到此這篇關(guān)于Java分布式學(xué)習(xí)之Kafka消息隊(duì)列的文章就介紹到這了,更多相關(guān)Java Kafka內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
解決springboot啟動(dòng)成功,但訪問(wèn)404的問(wèn)題
這篇文章主要介紹了解決springboot啟動(dòng)成功,但訪問(wèn)404的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-07-07
使用Spring Security和JWT實(shí)現(xiàn)安全認(rèn)證機(jī)制
在現(xiàn)代 Web 應(yīng)用中,安全認(rèn)證和授權(quán)是保障數(shù)據(jù)安全和用戶隱私的核心機(jī)制,Spring Security 是 Spring 框架下專為安全設(shè)計(jì)的模塊,具有高度的可配置性和擴(kuò)展性,而 JWT則是當(dāng)前流行的認(rèn)證解決方案,所以本文介紹了如何使用Spring Security和JWT實(shí)現(xiàn)安全認(rèn)證機(jī)制2024-11-11
Springboot如何實(shí)現(xiàn)Web系統(tǒng)License授權(quán)認(rèn)證
這篇文章主要介紹了Springboot如何實(shí)現(xiàn)Web系統(tǒng)License授權(quán)認(rèn)證,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-05-05
Springboot快速集成sse服務(wù)端推流(最新整理)
SSE?Server-Sent?Events是一種允許服務(wù)器向客戶端推送實(shí)時(shí)數(shù)據(jù)的技術(shù),它建立在?HTTP?和簡(jiǎn)單文本格式之上,提供了一種輕量級(jí)的服務(wù)器推送方式,通常也被稱為“事件流”(Event?Stream),這篇文章主要介紹了Springboot快速集成sse服務(wù)端推流(最新整理),需要的朋友可以參考下2024-02-02
java實(shí)現(xiàn)滑動(dòng)驗(yàn)證解鎖
這篇文章主要為大家詳細(xì)介紹了java實(shí)現(xiàn)滑動(dòng)驗(yàn)證解鎖,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2020-07-07
Spring?@bean和@component注解區(qū)別
本文主要介紹了Spring?@bean和@component注解區(qū)別,文中通過(guò)示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2022-01-01
在Idea2020.1中使用gitee2020.1.0創(chuàng)建第一個(gè)代碼庫(kù)的實(shí)現(xiàn)
這篇文章主要介紹了在Idea2020.1中使用gitee2020.1.0創(chuàng)建第一個(gè)代碼庫(kù)的實(shí)現(xiàn),文中通過(guò)圖文示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2020-07-07

