Springboot項(xiàng)目消費(fèi)Kafka數(shù)據(jù)的方法
一、引入依賴
你需要在 pom.xml 中添加 spring-kafka 相關(guān)依賴:
<dependencies>
<!-- Spring Boot Web -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- Spring Kafka -->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<!-- Spring Boot Starter for Logging (optional but useful for debugging) -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-logging</artifactId>
</dependency>
<!-- Spring Boot Starter for Testing -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>二、添加Kafka配置
在 application.yml 或 application.properties 文件中配置 Kafka 連接屬性:
application.yml 示例:
spring:
kafka:
bootstrap-servers: localhost:9092 # Kafka服務(wù)器地址
consumer:
group-id: my-consumer-group # 消費(fèi)者組ID
auto-offset-reset: earliest # 消費(fèi)者從頭開始讀?。ㄈ绻麤]有已提交的偏移量)
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer # 設(shè)置key的反序列化器
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer # 設(shè)置value的反序列化器為字符串
listener:
missing-topics-fatal: false # 如果主題不存在,不拋出致命錯(cuò)誤application.properties 示例:
spring.kafka.bootstrap-servers=localhost:9092 spring.kafka.consumer.group-id=my-consumer-group spring.kafka.consumer.auto-offset-reset=earliest spring.kafka.listener.missing-topics-fatal=false spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer # 設(shè)置key的反序列化器 spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer # 設(shè)置value的反序列化器為字符串
注意:spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer # 設(shè)置key的反序列化器
spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer # 設(shè)置value的反序列化器為字符串
以上配置說明Kafka生產(chǎn)的數(shù)據(jù)是json字符串,那么消費(fèi)接收的數(shù)據(jù)默認(rèn)也是json字符串,如果接收消息想用對(duì)象接受,需要自定義序列化器,比如以下配置
spring:
kafka:
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer # 對(duì) Key 使用 StringSerializer
value-serializer: org.springframework.kafka.support.serializer.ErrorHandlingSerializer # 對(duì) Value 使用 ErrorHandlingSerializer
properties:
spring.json.value.default.type: com.example.Order # 默認(rèn)的 JSON 反序列化目標(biāo)類型為 Order三、創(chuàng)建 Kafka 消費(fèi)者
創(chuàng)建一個(gè) Kafka 消費(fèi)者類來處理消息。你可以使用 @KafkaListener 注解來監(jiān)聽 Kafka 中的消息
(一)Kafka生產(chǎn)的消息是JSON 字符串
1、方式一
如果消息是 JSON 字符串,你可以使用 StringDeserializer 獲取消息后,再使用 ObjectMapper 將其轉(zhuǎn)換為
Java 對(duì)象(如 Order)。
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.stereotype.Service;
@Service
@EnableKafka // 啟用 Kafka 消費(fèi)者
public class KafkaConsumer {
private final ObjectMapper objectMapper = new ObjectMapper();
// 監(jiān)聽 Kafka 中的 order-topic 主題
@KafkaListener(topics = "order-topic", groupId = "order-consumer-group")
public void consumeOrder(String message) {
try {
// 將 JSON 字符串反序列化為 Order 對(duì)象
Order order = objectMapper.readValue(message, Order.class);
System.out.println("Received order: " + order);
} catch (Exception e) {
e.printStackTrace();
}
}
}說明:
@KafkaListener(topics = “my-topic”, groupId = “my-consumer-group”):
topics 表示監(jiān)聽的 Kafka 主題,groupId 表示消費(fèi)者所屬的消費(fèi)者組。
listen(String message): 該方法會(huì)被調(diào)用來處理收到的每條消息。在此示例中,我們打印出消息內(nèi)容。
2、方式二:需要直接訪問消息元數(shù)據(jù)
可以通過 ConsumerRecord 來接收 Kafka 消息。這種方式適用于需要直接訪問消息元數(shù)據(jù)(如
topic、partition、offset)的場(chǎng)景,也適合手動(dòng)管理消息消費(fèi)和偏移量提交的情況。
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;
@Service
public class KafkaConsumer {
// 監(jiān)聽 Kafka 中的 order-topic 主題
@KafkaListener(topics = "order-topic", groupId = "order-consumer-group")
public void consumeOrder(ConsumerRecord<String, String> record) {
// 獲取消息的詳細(xì)信息
String key = record.key(); // 獲取消息的 key
String value = record.value(); // 獲取消息的 value
String topic = record.topic(); // 獲取消息的 topic
int partition = record.partition(); // 獲取消息的分區(qū)
long offset = record.offset(); // 獲取消息的偏移量
long timestamp = record.timestamp(); // 獲取消息的時(shí)間戳
// 處理消息(這里我們只是打印消息)
System.out.println("Consumed record: ");
System.out.println("Key: " + key);
System.out.println("Value: " + value);
System.out.println("Topic: " + topic);
System.out.println("Partition: " + partition);
System.out.println("Offset: " + offset);
System.out.println("Timestamp: " + timestamp);
}
}(二)Kafka生產(chǎn)的消息是對(duì)象Order
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;
@Service
public class KafkaConsumer {
// 監(jiān)聽 Kafka 中的 order-topic 主題
@KafkaListener(topics = "order-topic", groupId = "order-consumer-group")
public void consumeOrder(ConsumerRecord<String, Order> record) {
// 獲取消息的詳細(xì)信息
String key = record.key(); // 獲取消息的 key
Order value = record.value(); // 獲取消息的 value
String topic = record.topic(); // 獲取消息的 topic
int partition = record.partition(); // 獲取消息的分區(qū)
long offset = record.offset(); // 獲取消息的偏移量
long timestamp = record.timestamp(); // 獲取消息的時(shí)間戳
// 處理消息(這里我們只是打印消息)
System.out.println("Consumed record: ");
System.out.println("Key: " + key);
System.out.println("Value: " + value);
System.out.println("Topic: " + topic);
System.out.println("Partition: " + partition);
System.out.println("Offset: " + offset);
System.out.println("Timestamp: " + timestamp);
}
}四、創(chuàng)建 啟動(dòng)類
確保你的 Spring Boot 啟動(dòng)類正確配置了 Spring Boot 應(yīng)用程序啟動(dòng)。
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class KafkaConsumerApplication {
public static void main(String[] args) {
SpringApplication.run(KafkaConsumerApplication.class, args);
}
}五、配置 Kafka 生產(chǎn)者(可選)
(一)消息類型為json串
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.stereotype.Service;
import com.fasterxml.jackson.databind.ObjectMapper;
@Service
@EnableKafka
public class KafkaProducer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate; // 發(fā)送的是 String 類型消息
private ObjectMapper objectMapper = new ObjectMapper(); // Jackson ObjectMapper 用于序列化
// 發(fā)送訂單到 Kafka
public void sendOrder(String topic, Order order) {
try {
// 將 Order 對(duì)象轉(zhuǎn)換為 JSON 字符串
String orderJson = objectMapper.writeValueAsString(order);
// 發(fā)送 JSON 字符串到 Kafka
kafkaTemplate.send(topic, orderJson); // 發(fā)送字符串消息
System.out.println("Order JSON sent to Kafka: " + orderJson);
} catch (Exception e) {
e.printStackTrace();
}
}
}(二)消息類型為對(duì)象Order
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.stereotype.Service;
@Service
@EnableKafka
public class KafkaProducer {
@Autowired
private KafkaTemplate<String, Order> kafkaTemplate;
// 發(fā)送訂單到 Kafka
public void sendOrder(String topic, Order order) {
kafkaTemplate.send(topic, order); // 發(fā)送訂單對(duì)象,Spring Kafka 會(huì)自動(dòng)將 Order 轉(zhuǎn)換為 JSON
}
}六、啟動(dòng) Kafka 服務(wù)
啟動(dòng) Kafka 服務(wù)
bin/kafka-server-start.sh config/server.properties
七、測(cè)試 Kafka 消費(fèi)者
你可以通過向 Kafka 發(fā)送消息來測(cè)試消費(fèi)者是否工作正常。假設(shè)你已經(jīng)在 Kafka 中創(chuàng)建了一個(gè)名為 my-topic 的主題,可以使用 KafkaProducer 來發(fā)送消息:
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class KafkaController {
@Autowired
private KafkaProducer kafkaProducer;
@GetMapping("/sendOrder")
public String sendOrder() {
Order order = new Order();
order.setOrderId(1L);
order.setUserId(123L);
order.setProduct("Laptop");
order.setQuantity(2);
order.setStatus("Created");
kafkaProducer.sendOrder("order-topic", order);
return "Order sent!";
}
}當(dāng)你訪問 /sendOrder端點(diǎn)時(shí),KafkaProducer 會(huì)將消息發(fā)送到 Kafka,KafkaConsumer 會(huì)接收到這條消息并打印出來。
九、測(cè)試和調(diào)試
你可以通過查看 Kafka 消費(fèi)者日志,確保消息已經(jīng)被成功消費(fèi)。你還可以使用 KafkaTemplate 發(fā)送消息,并確保 Kafka 生產(chǎn)者和消費(fèi)者之間的連接正常。
十、 結(jié)語
至此,你已經(jīng)在 Spring Boot 中成功配置并實(shí)現(xiàn)了 Kafka 消費(fèi)者和生產(chǎn)者。你可以根據(jù)需要擴(kuò)展功能,例如處理更復(fù)雜的消息類型、批量消費(fèi)等。
到此這篇關(guān)于Springboot項(xiàng)目如何消費(fèi)Kafka數(shù)據(jù)的文章就介紹到這了,更多相關(guān)Springboot消費(fèi)Kafka數(shù)據(jù)內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- springboot配置kafka批量消費(fèi),并發(fā)消費(fèi)方式
- SpringBoot整合Kafka完成生產(chǎn)消費(fèi)的方案
- springboot如何開啟和關(guān)閉kafka消費(fèi)
- spring-kafka使消費(fèi)者動(dòng)態(tài)訂閱新增的topic問題
- springboot集成kafka消費(fèi)手動(dòng)啟動(dòng)停止操作
- Springboot集成Kafka進(jìn)行批量消費(fèi)及踩坑點(diǎn)
- spring 整合kafka監(jiān)聽消費(fèi)的配置過程
- 淺談spring-kafka消費(fèi)異常處理
相關(guān)文章
詳解Java Optional正確使用方式和優(yōu)勢(shì)(避免空指針異常)
作為一個(gè) Java 后端開發(fā)者,NullPointerException(空指針異常)幾乎是我們寫代碼時(shí)最常見、最難纏的 Bug 之一,下面我們就來聊聊如何正確使用Optional以避免空指針異常吧2025-07-07
詳解Java如何通過Socket實(shí)現(xiàn)查詢IP
在本文中,我們來學(xué)習(xí)下如何找到連接到服務(wù)器的客戶端計(jì)算機(jī)的IP地址。我們將創(chuàng)建一個(gè)簡單的客戶端-服務(wù)器場(chǎng)景,讓我們探索用于TCP/IP通信的java.net?API,感興趣的可以了解一下2022-10-10
Java保留兩位小數(shù)的實(shí)現(xiàn)方法
這篇文章主要介紹了 Java保留兩位小數(shù)的實(shí)現(xiàn)方法的相關(guān)資料,需要的朋友可以參考下2017-06-06
解決SpringBoot jar包中的文件讀取問題實(shí)現(xiàn)
這篇文章主要介紹了解決SpringBoot jar包中的文件讀取問題實(shí)現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2020-08-08
如何解決java.lang.ClassNotFoundException: com.mysql.jdbc.Dr
這篇文章主要介紹了如何解決java.lang.ClassNotFoundException: com.mysql.jdbc.Driver問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-12-12
Java中l(wèi)ist.foreach()和list.stream().foreach()用法詳解
在Java中List是一種常用的集合類,用于存儲(chǔ)一組元素,List提供了多種遍歷元素的方式,包括使用forEach()方法和使用Stream流的forEach()方法,這篇文章主要給大家介紹了關(guān)于Java中l(wèi)ist.foreach()和list.stream().foreach()用法的相關(guān)資料,需要的朋友可以參考下2024-07-07
SpringCloud升級(jí)2020.0.x版之OpenFeign簡介與使用實(shí)現(xiàn)思路
在微服務(wù)系統(tǒng)中,我們經(jīng)常會(huì)進(jìn)行 RPC 調(diào)用。在 Spring Cloud 體系中,RPC 調(diào)用一般就是 HTTP 協(xié)議的調(diào)用。對(duì)于每次調(diào)用,都要經(jīng)過一系列詳細(xì)步驟,接下來通過本文給大家介紹SpringCloud OpenFeign簡介與使用,感興趣的朋友一起看看吧2021-10-10

