springboot使用kafka的過(guò)程
啟動(dòng)kafka
確保本地已安裝并啟動(dòng) Kafka 服務(wù)(或連接遠(yuǎn)程 Kafka 集群 ),比如通過(guò) Kafka 官網(wǎng)下載解壓后,啟動(dòng) Zookeeper(老版本 Kafka 依賴(lài),新版本用 KRaft 可不依賴(lài) )和 Kafka 服務(wù):
# 啟動(dòng) Zookeeper(若用 KRaft 模式可跳過(guò))
bin/zookeeper-server-start.sh config/zookeeper.properties
# 啟動(dòng) Kafka 服務(wù)
bin/kafka-server-start.sh config/server.properties
版本:
Kafka 從2.8.0版本開(kāi)始引入了 KIP-500,提供了無(wú) Zookeeper 的早期訪問(wèn)功能1。不過(guò),此時(shí)的實(shí)現(xiàn)并不完全,不建議在生產(chǎn)環(huán)境中使用。
3.0版本開(kāi)始真正全面摒棄 Zookeeper,使用新的元數(shù)據(jù)管理方式 Kraft,提高了 Kafka 的可擴(kuò)展性、可用性和性能4。
4.0版本是第一個(gè)完全無(wú)需 Apache Zookeeper 運(yùn)行的重大版本,將不再支持以 ZK 模式運(yùn)行或從 ZK 模式遷移。
項(xiàng)目引依賴(lài)
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.6.0</version> <!-- 版本按需選,建議用較新穩(wěn)定版 -->
</dependency>
</dependencies>創(chuàng)建 Producer 類(lèi)(編寫(xiě)生產(chǎn)者代碼)
import org.apache.kafka.clients.producer.*;
import java.util.Properties;
public class KafkaProducerDemo {
public static void main(String[] args) {
// 1. 配置 Kafka 連接、序列化等參數(shù)
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092"); // Kafka 集群地址
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); // 鍵的序列化器
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); // 值的序列化器
// 2. 創(chuàng)建 Producer 實(shí)例
Producer<String, String> producer = new KafkaProducer<>(props);
// 3. 構(gòu)造消息(指定主題、鍵、值)
String topic = "test_topic"; // 要發(fā)送到的主題,需提前在 Kafka 創(chuàng)建或允許自動(dòng)創(chuàng)建
String key = "key1";
String value = "Hello, Kafka from IDEA!";
ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value);
// 4. 發(fā)送消息(異步發(fā)送 + 回調(diào)處理結(jié)果)
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception exception) {
if (exception != null) {
System.err.println("消息發(fā)送失敗:" + exception.getMessage());
} else {
System.out.printf("消息發(fā)送成功!主題:%s,分區(qū):%d,偏移量:%d%n",
metadata.topic(), metadata.partition(), metadata.offset());
}
}
});
// 5. 關(guān)閉 Producer(實(shí)際生產(chǎn)環(huán)境可能在程序結(jié)束時(shí)或合適時(shí)機(jī)關(guān)閉)
producer.close();
}
}import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import java.util.Properties;
public class KafkaProducerExample {
private final static String TOPIC = "mytopic";
private final static String BOOTSTRAP_SERVERS = "localhost:9092";
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", BOOTSTRAP_SERVERS);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
try {
for (int i = 0; i < 10; i++) {
String message = "Message " + i;
producer.send(new ProducerRecord<>(TOPIC, message));
}
} catch (Exception e) {
e.printStackTrace();
} finally {
producer.close();
}
}
}創(chuàng)建 Consumer 類(lèi)(編寫(xiě)消費(fèi)者代碼)
import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class KafkaConsumerDemo {
public static void main(String[] args) {
// 1. 配置 Kafka 連接、反序列化、消費(fèi)者組等參數(shù)
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092"); // Kafka 集群地址
props.put("group.id", "test_group"); // 消費(fèi)者組 ID,同一組內(nèi)消費(fèi)者協(xié)調(diào)消費(fèi)
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); // 鍵的反序列化器
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); // 值的反序列化器
props.put("auto.offset.reset", "earliest"); // 沒(méi)有已提交偏移量時(shí),從最早消息開(kāi)始消費(fèi)
// 2. 創(chuàng)建 Consumer 實(shí)例
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
// 3. 訂閱主題
String topic = "test_topic";
consumer.subscribe(Collections.singletonList(topic));
// 4. 循環(huán)拉取消息(長(zhǎng)輪詢(xún))
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("收到消息:主題=%s,分區(qū)=%d,偏移量=%d,鍵=%s,值=%s%n",
record.topic(), record.partition(), record.offset(),
record.key(), record.value());
}
// 手動(dòng)提交偏移量(也可配置自動(dòng)提交,生產(chǎn)環(huán)境建議手動(dòng)更可靠)
consumer.commitSync();
}
} catch (Exception e) {
e.printStackTrace();
} finally {
// 5. 關(guān)閉 Consumer
consumer.close();
}
}
}import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.util.Collections;
import java.util.Properties;
public class KafkaConsumerExample {
private final static String TOPIC = "mytopic";
private final static String BOOTSTRAP_SERVERS = "localhost:9092";
private final static String GROUP_ID = "mygroup";
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", BOOTSTRAP_SERVERS);
props.put("group.id", GROUP_ID);
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList(TOPIC));
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(100);
// 處理接收到的消息
records.forEach(record -> {
System.out.println("Received message: " + record.value());
});
}
} catch (Exception e) {
e.printStackTrace();
} finally {
consumer.close();
}
}
}必須要素:
- 必要配置:
bootstrap.servers:Kafka 集群地址。group.id:消費(fèi)者組 ID(相同組內(nèi)的消費(fèi)者會(huì)負(fù)載均衡消費(fèi))。key.deserializer和value.deserializer:消息鍵和值的反序列化器。auto.offset.reset:消費(fèi)位置重置策略(如earliest從最早消息開(kāi)始消費(fèi))。
- 訂閱主題:通過(guò)
consumer.subscribe()訂閱目標(biāo)主題。 - 消息消費(fèi):通過(guò)
consumer.poll()輪詢(xún)拉取消息,并處理ConsumerRecords。 - 偏移量管理:自動(dòng)提交(
enable.auto.commit=true)或手動(dòng)提交(consumer.commitSync())消費(fèi)偏移量。 - 資源管理:使用后調(diào)用
consumer.close()關(guān)閉連接。
與 Kafka 的對(duì)比
Kafka的Producer和Consumer需要手動(dòng)管理連接和資源的關(guān)閉,因此在使用完畢后需要調(diào)用close方法來(lái)關(guān)閉Producer(或Consumer)。
總結(jié)來(lái)說(shuō),可以使用KafkaProducer的send方法來(lái)替代RabbitTemplate的convertAndSend方法在Kafka中發(fā)送消息。
Spring AMQP 是 Spring 框架提供的一個(gè)用于簡(jiǎn)化 AMQP(Advanced Message Queuing Protocol) 消息中間件開(kāi)發(fā)的模塊。它基于 AMQP 協(xié)議,提供了一套高層抽象和模板類(lèi),幫助開(kāi)發(fā)者更便捷地實(shí)現(xiàn)消息發(fā)送和接收,支持多種 AMQP 消息中間件(如 RabbitMQ、Apache Qpid 等)。
| 維度 | Spring AMQP(RabbitMQ) | Spring Kafka |
|---|---|---|
| 協(xié)議 | AMQP(高級(jí)消息隊(duì)列協(xié)議) | Kafka 自研協(xié)議 |
| 消息模型 | 支持多種交換器類(lèi)型(Direct、Topic 等) | 基于主題(Topic)和分區(qū)(Partition) |
| 順序性 | 單隊(duì)列內(nèi)保證順序 | 分區(qū)內(nèi)保證順序,多分區(qū)需按 Key 路由 |
| 吞吐量 | 中等(萬(wàn)級(jí) TPS) | 高(十萬(wàn)級(jí) TPS) |
| 適用場(chǎng)景 | 企業(yè)集成、任務(wù)調(diào)度、事務(wù)性消息 | 大數(shù)據(jù)、日志收集、實(shí)時(shí)流處理 |
到此這篇關(guān)于springboot使用kafka的文章就介紹到這了,更多相關(guān)springboot使用kafka內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Springboot的自動(dòng)配置是什么及注意事項(xiàng)
SpringBoot的自動(dòng)配置(Auto-configuration)是指框架根據(jù)項(xiàng)目的依賴(lài)和應(yīng)用程序的環(huán)境自動(dòng)配置Spring應(yīng)用上下文中的Bean和組件,目的是簡(jiǎn)化開(kāi)發(fā)者的配置工作,本文介紹Springboot的自動(dòng)配置是什么及注意事項(xiàng),感興趣的朋友一起看看吧2025-03-03
Java將對(duì)象保存到文件中/從文件中讀取對(duì)象的方法
下面小編就為大家?guī)?lái)一篇Java將對(duì)象保存到文件中/從文件中讀取對(duì)象的方法。小編覺(jué)得挺不錯(cuò)的,現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧2016-12-12
Mybatis-Plus自動(dòng)填充的實(shí)現(xiàn)示例
這篇文章主要介紹了Mybatis-Plus自動(dòng)填充的實(shí)現(xiàn)示例,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2019-08-08
Spring Boot 自動(dòng)配置的底層實(shí)現(xiàn)原理深度解析
SpringBoot自動(dòng)配置的核心在于通過(guò)“約定+條件判斷”實(shí)現(xiàn)配置的自動(dòng)化加載與生效,其底層實(shí)現(xiàn)可以拆解為觸發(fā)入口、配置類(lèi)加載、條件過(guò)濾、Bean注冊(cè)和配置覆蓋五個(gè)核心環(huán)節(jié),本文介紹Spring Boot自動(dòng)配置的底層實(shí)現(xiàn)原理,感興趣的朋友一起看看吧2025-12-12
解析Spring Boot內(nèi)嵌tomcat關(guān)于getServletContext().getRealPath獲取得到臨時(shí)
大家都很糾結(jié)這個(gè)問(wèn)題在使用getServletContext().getRealPath()得到的是臨時(shí)文件的路徑,每次重啟服務(wù),這個(gè)臨時(shí)文件的路徑還好變更,下面小編通過(guò)本文給大家分享Spring Boot內(nèi)嵌tomcat關(guān)于getServletContext().getRealPath獲取得到臨時(shí)路徑的問(wèn)題,一起看看吧2021-05-05

