最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Java中Springboot集成Kafka實現(xiàn)消息發(fā)送和接收功能

 更新時間:2025年01月24日 14:38:54   作者:cccl.  
Kafka是一個高吞吐量的分布式發(fā)布-訂閱消息系統(tǒng),主要用于處理大規(guī)模數(shù)據(jù)流,它由生產(chǎn)者、消費(fèi)者、主題、分區(qū)和代理等組件構(gòu)成,Kafka可以實現(xiàn)消息隊列、數(shù)據(jù)存儲和流處理等功能,在Java中,可以使用Spring Boot集成Kafka實現(xiàn)消息的發(fā)送和接收,感興趣的朋友跟隨小編一起看看吧

一、Kafka 簡介

Kafka 是由 Apache 軟件基金會開發(fā)的一個開源流處理平臺,最初由 LinkedIn 公司開發(fā),并于 2011 年開源。它是一種高吞吐量的分布式發(fā)布 - 訂閱消息系統(tǒng),以可持久化、高吞吐、低延遲、高容錯等特性而著稱。
Kafka 主要由生產(chǎn)者(Producer)、消費(fèi)者(Consumer)、主題(Topic)、分區(qū)(Partition)和代理(Broker)等組件構(gòu)成。生產(chǎn)者負(fù)責(zé)將數(shù)據(jù)發(fā)送到 Kafka 集群,消費(fèi)者從集群中讀取數(shù)據(jù)。主題是一種邏輯上的分類,數(shù)據(jù)被發(fā)送到特定的主題。每個主題又可以劃分為多個分區(qū),以實現(xiàn)數(shù)據(jù)的并行處理和提高系統(tǒng)的可擴(kuò)展性。代理則是 Kafka 集群中的服務(wù)器節(jié)點(diǎn),負(fù)責(zé)接收和存儲生產(chǎn)者發(fā)送的數(shù)據(jù),并為消費(fèi)者提供數(shù)據(jù)讀取服務(wù)。

二、Kafka 功能

消息隊列功能:Kafka 可以作為消息隊列使用,在應(yīng)用程序之間傳遞消息。生產(chǎn)者將消息發(fā)送到主題,不同的消費(fèi)者可以從主題中訂閱并消費(fèi)消息,實現(xiàn)應(yīng)用程序解耦。例如,在電商系統(tǒng)中,訂單生成模塊可以將訂單消息發(fā)送到 Kafka 主題,后續(xù)的庫存管理、物流配送等模塊可以從該主題消費(fèi)訂單消息,各自獨(dú)立處理,降低模塊間的耦合度。
數(shù)據(jù)存儲功能:Kafka 具有持久化存儲能力,它將消息數(shù)據(jù)存儲在磁盤上,并且通過多副本機(jī)制保證數(shù)據(jù)的可靠性。即使某個節(jié)點(diǎn)出現(xiàn)故障,數(shù)據(jù)也不會丟失。這種特性使得 Kafka 不僅可以作為消息隊列,還能用于數(shù)據(jù)的長期存儲和備份,例如用于存儲系統(tǒng)的操作日志,方便后續(xù)的數(shù)據(jù)分析和故障排查。
流處理功能:Kafka 可以與流處理框架(如 Apache Flink、Spark Streaming 等)集成,對實時數(shù)據(jù)流進(jìn)行處理。通過將實時數(shù)據(jù)發(fā)送到 Kafka 主題,流處理框架可以從主題中讀取數(shù)據(jù)并進(jìn)行實時計算、分析和轉(zhuǎn)換。例如,在實時監(jiān)控系統(tǒng)中,通過 Kafka 收集服務(wù)器的性能指標(biāo)數(shù)據(jù),然后使用流處理框架對這些數(shù)據(jù)進(jìn)行實時分析,及時發(fā)現(xiàn)性能異常并發(fā)出警報。

三、POM依賴

    <!-- kafka-->
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
        <version>2.8.11</version>
    </dependency>

四、配置文件

spring:
  # Kafka 配置
  kafka:
    # Kafka 服務(wù)器地址和端口 代理地址,可以多個
    bootstrap-servers: IP:9092
    # 生產(chǎn)者配置
    producer:
      # 發(fā)送失敗時的重試次數(shù)
      retries: 3
      # 每次批量發(fā)送消息的數(shù)量,調(diào)整為較小值
      batch-size: 1
      # 生產(chǎn)者緩沖區(qū)大小
      buffer-memory: 33554432
      # 消息 key 的序列化器,將 key 序列化為字節(jié)數(shù)組
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      # 消息 value 的序列化器,將消息體序列化為字節(jié)數(shù)組
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
    # 消費(fèi)者配置
    consumer:
      # 當(dāng)沒有初始偏移量或當(dāng)前偏移量不存在時,從最早的消息開始消費(fèi)
      auto-offset-reset: earliest
      # 是否自動提交偏移量
      enable-auto-commit: true
      # 自動提交偏移量的時間間隔(毫秒),延長自動提交時間間隔
      auto-commit-interval: 1000
      # 消息 key 的反序列化器,將字節(jié)數(shù)組反序列化為 key
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      # 消息 value 的反序列化器,將字節(jié)數(shù)組反序列化為消息體
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer

五、生產(chǎn)者

import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Component;
import org.springframework.util.concurrent.ListenableFuture;
import org.springframework.util.concurrent.ListenableFutureCallback;
/**
 * 生產(chǎn)者
 *
 * @author chenlei
 */
@Slf4j
@Component
public class KafkaProducer {
    /**
     * KafkaTemplate
     */
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;
    /**
     * 發(fā)送消息到指定的 Kafka 主題,并可指定分組信息
     *
     * @param topic   消息要發(fā)送到的 Kafka 主題
     * @param message 要發(fā)送的消息內(nèi)容
     */
    public void sendMessage(String topic, String message) {
        // 使用 KafkaTemplate 發(fā)送消息,將消息發(fā)送到指定的主題
        ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
        future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
            @Override
            public void onSuccess(SendResult<String, String> result) {
                // 消息發(fā)送成功后的處理邏輯,可根據(jù)需要添加
                log.info("已發(fā)送消息=[" + message + "],其偏移量=[" + result.getRecordMetadata().offset() + "]");
            }
            @Override
            public void onFailure(Throwable ex) {
                // 消息發(fā)送失敗后的處理邏輯,使用日志記錄異常
                log.error("發(fā)送消息=[" + message + "] 失敗", ex);
            }
        });
    }
}

六、消費(fèi)者

import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
/**
 * @author 消費(fèi)者
 * chenlei
 */
@Slf4j
@Component
public class KafkaConsumer {
    /**
     * 監(jiān)聽 Kafka 主題方法。
     *
     * @param record 從 Kafka 接收到的 ConsumerRecord,包含消息的鍵值對
     */
    @KafkaListener(topics = {"topic"}, groupId = "consumer.group-id", concurrency = "5")
    public void listen(ConsumerRecord<?, ?> record) {
        // 打印接收到的消息的詳細(xì)信息
        log.info("接收到 Kafka 消息: 主題 = {}, 分區(qū) = {}, 偏移量 = {}, 鍵 = {}, 值 = {}",
                record.topic(), record.partition(), record.offset(), record.key(), record.value());
    }
}

到此這篇關(guān)于Java中Springboot集成Kafka實現(xiàn)消息發(fā)送和接收的文章就介紹到這了,更多相關(guān)Springboot Kafka 消息發(fā)送和接收內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 解決java頁面URL地址傳輸參數(shù)亂碼的方法

    解決java頁面URL地址傳輸參數(shù)亂碼的方法

    這篇文章主要介紹了解決java頁面URL地址傳輸參數(shù)亂碼的方法,URL地址參數(shù)亂碼問題,算是老話重談了吧!需要的朋友可以參考下
    2015-09-09
  • java使用stream判斷兩個list元素的屬性并輸出方式

    java使用stream判斷兩個list元素的屬性并輸出方式

    這篇文章主要介紹了java使用stream判斷兩個list元素的屬性并輸出方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • Java import導(dǎo)入及訪問控制權(quán)限修飾符原理解析

    Java import導(dǎo)入及訪問控制權(quán)限修飾符原理解析

    這篇文章主要介紹了Java import導(dǎo)入及訪問控制權(quán)限修飾符過程解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2019-11-11
  • 詳解Java中自定義注解的使用

    詳解Java中自定義注解的使用

    Annontation是Java5開始引入的新特征,中文名稱叫注解,它提供了一種安全的類似注釋的機(jī)制,用來將任何的信息或元數(shù)據(jù)(metadata)與程序元素(類、方法、成員變量等)進(jìn)行關(guān)聯(lián)。本文主要介紹了自定義注解的使用,希望對大家有所幫助
    2023-03-03
  • Java Swing窗體關(guān)閉事件的調(diào)用關(guān)系

    Java Swing窗體關(guān)閉事件的調(diào)用關(guān)系

    這篇文章主要為大家詳細(xì)介紹了Java Swing窗體關(guān)閉事件的調(diào)用關(guān)系,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2019-07-07
  • Jsoup解析HTML實例及文檔方法詳解

    Jsoup解析HTML實例及文檔方法詳解

    這篇文章主要介紹了Jsoup如何解析一個HTML文檔、從文件加載文檔、從URL加載Document等方法,對Jsoup常用方法做了詳細(xì)講解,最近提供了一個示例供大家參考 使用DOM方法來遍歷一個文檔 從元素抽取屬性,文本和HTML 獲取所有鏈接
    2013-11-11
  • Java中Integer.equals的用法與特殊情況

    Java中Integer.equals的用法與特殊情況

    Java中Integer.equals比較數(shù)值而非對象地址,自動裝箱處理int參數(shù),-128~127范圍內(nèi)使用==安全,超出范圍或null需用Objects.equals,compareTo用于大小比較,equals僅判斷相等性,注意類型不匹配會導(dǎo)致false
    2025-07-07
  • Spring整合Mybatis方式之注冊映射器

    Spring整合Mybatis方式之注冊映射器

    這篇文章主要介紹了Spring整合Mybatis方式之注冊映射器,MapperFactoryBean注冊映射器的最大問題,就是需要一個個注冊所有的映射器,而實際上mybatis-spring提供了掃描包下所有映射器接口的方法,每種方式給大家介紹的非常詳細(xì),需要的朋友參考下吧
    2024-03-03
  • java實現(xiàn)簡單控制臺通訊錄

    java實現(xiàn)簡單控制臺通訊錄

    這篇文章主要為大家詳細(xì)介紹了java實現(xiàn)簡單控制臺通訊錄,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2018-02-02
  • Java事件機(jī)制要素及實例詳解

    Java事件機(jī)制要素及實例詳解

    這篇文章主要介紹了Java事件機(jī)制要素及實例詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2020-04-04

最新評論

新密市| 庆云县| 齐河县| 拉萨市| 乾安县| 蒲江县| 泌阳县| 麻江县| 东宁县| 固始县| 德保县| 祥云县| 沿河| 宜兰市| 襄樊市| 靖西县| 枣庄市| 吴川市| 仁化县| 陇南市| 高台县| 香港| 青岛市| 茌平县| 芷江| 临汾市| 页游| 贡嘎县| 山丹县| 永顺县| 泰和县| 福清市| 三门县| 宁海县| 泰顺县| 汕头市| 太白县| 钟祥市| 凯里市| 天气| 广汉市|