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

springboot使用kafka的過(guò)程

 更新時(shí)間:2025年06月18日 09:17:31   作者:小魚(yú)小魚(yú).oO  
本文介紹了Spring Boot集成Kafka的步驟,包括啟動(dòng)服務(wù)、配置生產(chǎn)者與消費(fèi)者,以及Kafka從依賴(lài)Zookeeper到Kraft模式的版本演進(jìn),本文結(jié)合實(shí)例代碼給大家介紹的非常詳細(xì),需要的朋友參考下吧

啟動(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)配置是什么及注意事項(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設(shè)計(jì)模式之享元模式示例詳解

    Java設(shè)計(jì)模式之享元模式示例詳解

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

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

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

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

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

    Java JVM類(lèi)加載機(jī)制解讀

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

    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)原理深度解析

    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
  • 隱藏idea的.idea和.mvn文件的解決方案

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

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

    Java私有構(gòu)造器使用方法示例

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

    解析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

最新評(píng)論

当涂县| 商水县| 吉水县| 巴彦淖尔市| 靖安县| 双峰县| 桃园县| 即墨市| 台江县| 格尔木市| 祁门县| 平南县| 泽普县| 禄劝| 洱源县| 芮城县| 宁波市| 凤翔县| 东丽区| 任丘市| 鄂托克前旗| 上栗县| 明光市| 旬阳县| 达日县| 永昌县| 上犹县| 吴江市| 安徽省| 天水市| 措美县| 酒泉市| 克东县| 呼伦贝尔市| 乐山市| 吉安市| 康马县| 永年县| 绥江县| 青阳县| 科技|