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

springboot使用kafka的過程

 更新時間:2025年06月18日 09:17:31   作者:小魚小魚.oO  
本文介紹了Spring Boot集成Kafka的步驟,包括啟動服務、配置生產者與消費者,以及Kafka從依賴Zookeeper到Kraft模式的版本演進,本文結合實例代碼給大家介紹的非常詳細,需要的朋友參考下吧

啟動kafka

確保本地已安裝并啟動 Kafka 服務(或連接遠程 Kafka 集群 ),比如通過 Kafka 官網下載解壓后,啟動 Zookeeper(老版本 Kafka 依賴,新版本用 KRaft 可不依賴 )和 Kafka 服務:

# 啟動 Zookeeper(若用 KRaft 模式可跳過)

bin/zookeeper-server-start.sh config/zookeeper.properties

# 啟動 Kafka 服務

bin/kafka-server-start.sh config/server.properties

版本:

Kafka 從2.8.0版本開始引入了 KIP-500,提供了無 Zookeeper 的早期訪問功能1。不過,此時的實現(xiàn)并不完全,不建議在生產環(huán)境中使用。

3.0版本開始真正全面摒棄 Zookeeper,使用新的元數(shù)據(jù)管理方式 Kraft,提高了 Kafka 的可擴展性、可用性和性能4。

4.0版本是第一個完全無需 Apache Zookeeper 運行的重大版本,將不再支持以 ZK 模式運行或從 ZK 模式遷移。

項目引依賴

<dependencies>
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>3.6.0</version> <!-- 版本按需選,建議用較新穩(wěn)定版 -->
    </dependency>
</dependencies>

創(chuàng)建 Producer 類(編寫生產者代碼)

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 實例
        Producer<String, String> producer = new KafkaProducer<>(props);
        // 3. 構造消息(指定主題、鍵、值)
        String topic = "test_topic"; // 要發(fā)送到的主題,需提前在 Kafka 創(chuàng)建或允許自動創(chuàng)建
        String key = "key1";
        String value = "Hello, Kafka from IDEA!";
        ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value);
        // 4. 發(fā)送消息(異步發(fā)送 + 回調處理結果)
        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. 關閉 Producer(實際生產環(huá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 類(編寫消費者代碼)

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 連接、反序列化、消費者組等參數(shù)
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092"); // Kafka 集群地址
        props.put("group.id", "test_group"); // 消費者組 ID,同一組內消費者協(xié)調消費
        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"); // 沒有已提交偏移量時,從最早消息開始消費
        // 2. 創(chuàng)建 Consumer 實例
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
        // 3. 訂閱主題
        String topic = "test_topic";
        consumer.subscribe(Collections.singletonList(topic));
        // 4. 循環(huá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());
                }
                // 手動提交偏移量(也可配置自動提交,生產環(huán)境建議手動更可靠)
                consumer.commitSync();
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            // 5. 關閉 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:消費者組 ID(相同組內的消費者會負載均衡消費)。
    • key.deserializer 和 value.deserializer:消息鍵和值的反序列化器。
    • auto.offset.reset:消費位置重置策略(如 earliest 從最早消息開始消費)。
  • 訂閱主題:通過 consumer.subscribe() 訂閱目標主題。
  • 消息消費:通過 consumer.poll() 輪詢拉取消息,并處理 ConsumerRecords。
  • 偏移量管理:自動提交(enable.auto.commit=true)或手動提交(consumer.commitSync())消費偏移量。
  • 資源管理:使用后調用 consumer.close() 關閉連接。

與 Kafka 的對比

Kafka的Producer和Consumer需要手動管理連接和資源的關閉,因此在使用完畢后需要調用close方法來關閉Producer(或Consumer)。

總結來說,可以使用KafkaProducer的send方法來替代RabbitTemplate的convertAndSend方法在Kafka中發(fā)送消息。

Spring AMQP 是 Spring 框架提供的一個用于簡化 AMQP(Advanced Message Queuing Protocol) 消息中間件開發(fā)的模塊。它基于 AMQP 協(xié)議,提供了一套高層抽象和模板類,幫助開發(fā)者更便捷地實現(xiàn)消息發(fā)送和接收,支持多種 AMQP 消息中間件(如 RabbitMQ、Apache Qpid 等)。

維度Spring AMQP(RabbitMQ)Spring Kafka
協(xié)議AMQP(高級消息隊列協(xié)議)Kafka 自研協(xié)議
消息模型支持多種交換器類型(Direct、Topic 等)基于主題(Topic)和分區(qū)(Partition)
順序性單隊列內保證順序分區(qū)內保證順序,多分區(qū)需按 Key 路由
吞吐量中等(萬級 TPS)高(十萬級 TPS)
適用場景企業(yè)集成、任務調度、事務性消息大數(shù)據(jù)、日志收集、實時流處理

到此這篇關于springboot使用kafka的文章就介紹到這了,更多相關springboot使用kafka內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • Springboot的自動配置是什么及注意事項

    Springboot的自動配置是什么及注意事項

    SpringBoot的自動配置(Auto-configuration)是指框架根據(jù)項目的依賴和應用程序的環(huán)境自動配置Spring應用上下文中的Bean和組件,目的是簡化開發(fā)者的配置工作,本文介紹Springboot的自動配置是什么及注意事項,感興趣的朋友一起看看吧
    2025-03-03
  • Java設計模式之享元模式示例詳解

    Java設計模式之享元模式示例詳解

    享元模式(FlyWeight?Pattern),也叫蠅量模式,運用共享技術,有效的支持大量細粒度的對象,享元模式就是池技術的重要實現(xiàn)方式。本文將通過示例詳細講解享元模式,感興趣的可以了解一下
    2022-03-03
  • Java文件上傳下載、郵件收發(fā)實例代碼

    Java文件上傳下載、郵件收發(fā)實例代碼

    這篇文章主要介紹了Java文件上傳下載、郵件收發(fā)實例代碼的相關資料,非常不錯具有參考借鑒價值,需要的朋友可以參考下
    2016-06-06
  • Java將對象保存到文件中/從文件中讀取對象的方法

    Java將對象保存到文件中/從文件中讀取對象的方法

    下面小編就為大家?guī)硪黄狫ava將對象保存到文件中/從文件中讀取對象的方法。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2016-12-12
  • Java JVM類加載機制解讀

    Java JVM類加載機制解讀

    JVM將class文件字節(jié)碼文件加載到內存中, 并將這些靜態(tài)數(shù)據(jù)轉換成方法區(qū)中的運行時數(shù)據(jù)結構,在堆(并不一定在堆中,HotSpot在方法區(qū)中)中生成一個代表這個類的java.lang.Class 對象,作為方法區(qū)類數(shù)據(jù)的訪問入口,接下來將詳細講解JVM類加載機制
    2021-11-11
  • Mybatis-Plus自動填充的實現(xiàn)示例

    Mybatis-Plus自動填充的實現(xiàn)示例

    這篇文章主要介紹了Mybatis-Plus自動填充的實現(xiàn)示例,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2019-08-08
  • Spring Boot 自動配置的底層實現(xiàn)原理深度解析

    Spring Boot 自動配置的底層實現(xiàn)原理深度解析

    SpringBoot自動配置的核心在于通過“約定+條件判斷”實現(xiàn)配置的自動化加載與生效,其底層實現(xiàn)可以拆解為觸發(fā)入口、配置類加載、條件過濾、Bean注冊和配置覆蓋五個核心環(huán)節(jié),本文介紹Spring Boot自動配置的底層實現(xiàn)原理,感興趣的朋友一起看看吧
    2025-12-12
  • 隱藏idea的.idea和.mvn文件的解決方案

    隱藏idea的.idea和.mvn文件的解決方案

    這篇文章主要介紹了隱藏idea的.idea和.mvn文件的解決方法,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2023-07-07
  • Java私有構造器使用方法示例

    Java私有構造器使用方法示例

    這篇文章主要介紹了Java私有構造器的含義、關鍵字,同時通過實例向大家展示其使用方法,需要的朋友可以參考下
    2017-09-09
  • 解析Spring Boot內嵌tomcat關于getServletContext().getRealPath獲取得到臨時路徑的問題

    解析Spring Boot內嵌tomcat關于getServletContext().getRealPath獲取得到臨時

    大家都很糾結這個問題在使用getServletContext().getRealPath()得到的是臨時文件的路徑,每次重啟服務,這個臨時文件的路徑還好變更,下面小編通過本文給大家分享Spring Boot內嵌tomcat關于getServletContext().getRealPath獲取得到臨時路徑的問題,一起看看吧
    2021-05-05

最新評論

丰原市| 沙洋县| 武川县| 江川县| 综艺| 上高县| 枣庄市| 襄汾县| 武强县| 澄迈县| 库伦旗| 安塞县| 江孜县| 松滋市| 曲靖市| 策勒县| 永胜县| 莱阳市| 遂平县| 县级市| 威海市| 山阳县| 曲靖市| 台湾省| 新昌县| 藁城市| 福鼎市| 闽清县| 天津市| 台中市| 临沭县| 长岭县| 从江县| 彰化市| 济阳县| 师宗县| 隆化县| 绥化市| 涪陵区| 思南县| 习水县|