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

Kafka?在?java?中的基本使用步驟

 更新時間:2025年09月05日 10:55:38   作者:張小虎在學(xué)習(xí)  
本文給大家介紹Kafka在java中的基本使用,本文分步驟結(jié)合實例代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友參考下吧

1. 使用 kafka 原生客戶端

現(xiàn)在基本都直接使用 springboot 版本,但了解原生客戶端,能更好的理解 springboot 版的 kafka 客戶端原理。

步驟1:pom 引入核心依賴:

引入依賴時,盡量選擇和 kafka 版本對應(yīng)的依賴版本。

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka_2.13</artifactId>
    <version>4.0.0</version>
</dependency>

步驟2:提供者客戶端代碼:

提供者客戶端要做三件事:

  1. 設(shè)置提供者客戶端屬性(可選屬性都被定義在 ProducerConfig 類中)
  2. 設(shè)置要發(fā)送的消息
  3. 發(fā)送(有三種發(fā)送方式,下面代碼中都有)
public class MyProducer {
    public static void main(String[] args) throws ExecutionException, InterruptedException {
        // 第一步:設(shè)置提供者屬性
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.2.28:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        try (Producer<String, String> producer = new KafkaProducer<>(props)) {
            // 第二步:設(shè)置要發(fā)送的消息
            ProducerRecord<String, String> record = new ProducerRecord<>("testTopic", "testKey", "testValue");
            // 第三部:發(fā)送消息
            // send(producer, record);
            // sendSync(producer, record);
            sendASync(producer, record);
        }
    }
    /**
     * 發(fā)送方式1:單向推送,不關(guān)心服務(wù)器的應(yīng)答
     */
    private static void send(Producer<String, String> producer, ProducerRecord<String, String> record) {
        producer.send(record);
    }
    /**
     * 發(fā)送方式2:同步推送,得到服務(wù)器的應(yīng)答前會阻塞當(dāng)前線程
     */
    private static void sendSync(Producer<String, String> producer, ProducerRecord<String, String> record) throws ExecutionException, InterruptedException {
        RecordMetadata metadata = producer.send(record).get();
        System.out.println(metadata.topic());
        System.out.println(metadata.partition());
        System.out.println(metadata.offset());
    }
    /**
     * 發(fā)送方式3:異步推送,不需等待服務(wù)器應(yīng)答,當(dāng)服務(wù)器有應(yīng)答后會觸發(fā)函數(shù)回調(diào)
     */
    private static void sendASync(Producer<String, String> producer, ProducerRecord<String, String> record) {
        producer.send(record, (metadata, exception) -> {
            if (exception != null) {
                throw new RuntimeException("向 kafka 推送失敗", exception);
            }
            System.out.println(metadata.topic());
            System.out.println(metadata.partition());
            System.out.println(metadata.offset());
        });
    }
}

步驟3:消費(fèi)者客戶端代碼:

消費(fèi)者客戶端要做三件事:

  1. 設(shè)置消費(fèi)者客戶端屬性(可選屬性都被定義在 ConsumerConfig 類中)
  2. 設(shè)置消費(fèi)者訂閱的主題
  3. 拉取消息
  4. 提交 offset(有兩種提交方式,下面代碼中都有)
public class MyConsumer {
    public static void main(String[] args) {
        // 第一步:設(shè)置消費(fèi)者屬性
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.2.28:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "testGroup");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        try (Consumer<String, String> consumer = new KafkaConsumer<>(props)) {
            // 第二步:設(shè)置要訂閱的主題
            consumer.subscribe(Collections.singletonList("testTopic"));
            while (true) {
                // 第三步:拉取消息,100 代表最大等待時間,如果時間到了還沒有拉取到消息就不阻塞了繼續(xù)往后執(zhí)行
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofNanos(100));
                for (ConsumerRecord<String, String> record : records) {
                    System.out.println(record.value());
                }
                // 第四步:提交 offset
                // consumer.commitSync(); // 同步提交,表示必須等到 offset 提交完畢,再去消費(fèi)下?批數(shù)據(jù)
                consumer.commitSync(); // 異步提交,表示發(fā)送完提交 offset 請求后,就開始消費(fèi)下?批數(shù)據(jù)了。不?等到Broker的確認(rèn)。
            }
        }
    }
}

2. Kafka 集成 springboot

springboot 版本是最常用的,比原生客戶端使用方便。但是道理是一樣的,底層也是原生客戶端。

pom 引入核心依賴:

<parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <version>3.1.0</version>
</parent>
<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>
</dependencies>

yaml 配置文件:

這一步無非就是把原生客戶端中的屬性配置,寫在 yaml 中

spring:
  kafka:
    bootstrap-servers: 192.168.2.28:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
    consumer:
      group-id: testGroup
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer

提供者客戶端代碼:

只需要兩步:

  1. 注入 KafkaTemplate
  2. 發(fā)送
@RestController
public class ProducerController {
    /**
     * kafka
     */
    private KafkaTemplate<String, Object> kafkaTemplate;
    @Autowired
    public void setKafkaTemplate(KafkaTemplate<String, Object> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }
    @GetMapping("/test")
    public void send() {
        // 發(fā)送 kafka 消息
        kafkaTemplate.send("testTopic", "testKey", "testValue");
    }
}

消費(fèi)者客戶端代碼:

只需要監(jiān)聽主題就可以

@RestController
public class ConsumerController {
    // 監(jiān)聽 kafka 消息
    @KafkaListener(topics = {"testTopic"})
    public void test(ConsumerRecord<?, ?> record) {
        System.out.println(record.value());
    }
}

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

相關(guān)文章

  • Java實現(xiàn)選擇排序算法的實例教程

    Java實現(xiàn)選擇排序算法的實例教程

    這篇文章主要介紹了Java實現(xiàn)選擇排序算法的實例教程,選擇排序的時間復(fù)雜度為О(n&sup2;),需要的朋友可以參考下
    2016-05-05
  • SpringBoot和前端Vue的跨域問題及解決方案

    SpringBoot和前端Vue的跨域問題及解決方案

    所謂跨域就是從 A 向 B 發(fā)請求,如若他們的地址協(xié)議、域名、端口都不相同,直接訪問就會造成跨域問題,跨域是非常常見的現(xiàn)象,這篇文章主要介紹了解決SpringBoot和前端Vue的跨域問題,需要的朋友可以參考下
    2023-11-11
  • Java中Controller引起的Ambiguous?mapping問題及解決

    Java中Controller引起的Ambiguous?mapping問題及解決

    這篇文章主要介紹了Java中Controller引起的Ambiguous?mapping問題及解決,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-10-10
  • Java通俗易懂系列設(shè)計模式之責(zé)任鏈模式

    Java通俗易懂系列設(shè)計模式之責(zé)任鏈模式

    這篇文章主要介紹了Java通俗易懂系列設(shè)計模式之責(zé)任鏈模式,對設(shè)計模式感興趣的同學(xué),一定要看一下
    2021-04-04
  • Spring Boot多模塊化后,服務(wù)間調(diào)用的坑及解決

    Spring Boot多模塊化后,服務(wù)間調(diào)用的坑及解決

    這篇文章主要介紹了Spring Boot多模塊化后,服務(wù)間調(diào)用的坑及解決,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-06-06
  • springboot整合druid的完整流程記錄

    springboot整合druid的完整流程記錄

    Java程序很大一部分要操作數(shù)據(jù)庫,為了提高性能操作數(shù)據(jù)庫的時候,又不得不使用數(shù)據(jù)庫連接池,Druid是阿里巴巴開源平臺上一個數(shù)據(jù)庫連接池實現(xiàn),結(jié)合了DBCP等DB池的優(yōu)點(diǎn),同時加入了日志監(jiān)控,這篇文章主要介紹了springboot整合druid的相關(guān)資料,需要的朋友可以參考下
    2026-02-02
  • Java實現(xiàn)簡單的模板渲染

    Java實現(xiàn)簡單的模板渲染

    這篇文章主要為大家詳細(xì)介紹了Java實現(xiàn)簡單的模板渲染的相關(guān)資料,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-12-12
  • 一文搞懂Java并發(fā)AQS的共享鎖模式

    一文搞懂Java并發(fā)AQS的共享鎖模式

    這篇文章主要為大家闡述AQS另外一個重要模式,共享鎖模式。共享鎖可以由多個線程同時獲取,?比較典型的就是讀鎖,感興趣的小伙伴可以了解一下
    2022-10-10
  • 多線程_解決Runnable接口無start()方法的情況

    多線程_解決Runnable接口無start()方法的情況

    這篇文章主要介紹了多線程_解決Runnable接口無start()方法的情況,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2021-03-03
  • SpringCloud Alibaba Sentinel 從入門到精通

    SpringCloud Alibaba Sentinel 從入門到精通

    本文介紹了Sentinel的核心概念、使用方式及規(guī)則配置,Sentinel是阿里巴巴開源的流量控制框架,主要應(yīng)用于流量控制、熔斷降級等場景,文中有詳細(xì)的規(guī)則配置示例,并介紹了如何與Feign整合,感興趣的朋友一起看看吧
    2026-04-04

最新評論

南川市| 新源县| 南雄市| 湖北省| 云南省| 衡水市| 长宁区| 剑河县| 太湖县| 农安县| 新干县| 皮山县| 明光市| 宁武县| 梅河口市| 鹿泉市| 江永县| 青浦区| 凤阳县| 探索| 绥江县| 福鼎市| 射阳县| 泾源县| 寻乌县| 阿巴嘎旗| 冷水江市| 五常市| 上杭县| 佛坪县| 绩溪县| 扎囊县| 马山县| 东安县| 马山县| 柘城县| 汝阳县| 宝山区| 兴隆县| 施秉县| 女性|