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

Springboot使用kafka的兩種方式

 更新時間:2023年11月08日 09:02:15   作者:香菜菜  
在公司用kafka比較多,今天整理下Springboot使用kafka的兩種方式,Kafka作為一個消息發(fā)布訂閱系統(tǒng),就包括消息生成者和消息消費者,文中通過代碼示例介紹的非常詳細,具有一定的參考價值,需要的朋友可以參考下

1、創(chuàng)建實驗項目

第一步創(chuàng)建一個Springboot項目,引入spring-kafka依賴,這是后面的基礎。

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka</artifactId>
        </dependency>

kafka配置

spring:
  kafka:
    bootstrap-servers: kafka.tyjt.com:9092
    consumer:
      auto-offset-reset: earliest
      group-id: sharingan-group
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.ByteArrayDeserializer

2、自動檔

為了方便使用kafka,Springboot提供了spring-kafka 這個包,在已開始我們已經(jīng)導入了,下面直接使用吧

Spring項目里引入Kafka非常方便,使用kafkaTemplate(Producer的模版)+@KafkaListener(Consumer的監(jiān)聽器)即可完成生產(chǎn)者-消費者的代碼開發(fā)

2.1 監(jiān)聽listener

為了使創(chuàng)建 kafka 監(jiān)聽器更加簡單,Spring For Kafka 提供了 @KafkaListener 注解,

@KafkaListener 注解配置方法上,凡是此注解的方法就會被標記為是 Kafka 消息監(jiān)聽器,所以可以用

@KafkaListener 注解快速創(chuàng)建消息監(jiān)聽器。

@Configuration
@EnableKafka
public class ConsumerConfigDemo {
    @KafkaListener(topics = {"test"},groupId = "group1")
    public void kafkaListener(String topic,String message){
        System.out.println("消息:"+message);
    }
}

2.2 發(fā)布消息

發(fā)布消息通過kafkaTemplate,kafkaTemplate是spring-kafka 的封裝

@Slf4j
@Service
public class KafkaProducerService {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;
    public void sendMessage(String topic, String key, String message) throws Exception {
        kafkaTemplate.send(topic,key,message);
    }
}

kafkaTemplate 有很多不同的發(fā)送方法,根據(jù)自己的需求使用,這里只記錄最簡單的狀況。

3、手動檔

3.1 手動創(chuàng)建consumer

關于consumer的主要的封裝在ConcurrentKafkaListenerContainerFactory這個里頭,

本身的KafkaConsumer是線程不安全的,無法并發(fā)操作,這里spring又在包裝了一層,根據(jù)配置的spring.kafka.listener.concurrency來生成多個并發(fā)的KafkaMessageListenerContainer實例

package com.tyjt.sharingan.kafka;
 
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.listener.AcknowledgingConsumerAwareMessageListener;
import org.springframework.kafka.listener.ConcurrentMessageListenerContainer;
import org.springframework.kafka.listener.ContainerProperties;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;
 
import javax.annotation.Resource;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
 
/**
 * 啟動kafka consumer
 *
 * @author 種鑫
 * @date 2023/10/18 17:26
 */
@EnableKafka
@Component
@Slf4j
public class KafkaConsumerMgr {
    @Resource
    ConcurrentKafkaListenerContainerFactory<String, byte[]> containerFactory;
 
    Map<String, ConcurrentMessageListenerContainer<?, ?>> containerMap = new ConcurrentHashMap<>();
    public void startListener(KafkaProtoConsumer kafkaConsumer) {
        //  停止相同的
        if (containerMap.containsKey(kafkaConsumer.getTopic())) {
            containerMap.get(kafkaConsumer.getTopic()).stop();
        }
        ConcurrentMessageListenerContainer<String, byte[]> container = createListenerContainer(kafkaConsumer);
        container.start();
        containerMap.put(kafkaConsumer.getTopic(), container);
    }
 
    private ConcurrentMessageListenerContainer<String, byte[]> createListenerContainer(KafkaProtoConsumer consumer) {
        ConcurrentMessageListenerContainer<String, byte[]> container = containerFactory.createContainer(consumer.topic());
        container.setBeanName(consumer.group() + "-" + consumer.topic());
        container.setConcurrency(consumer.getPartitionCount());
        consumer.deployContainer(container);
        // 防止被修改的配置
        ContainerProperties containerProperties = container.getContainerProperties();
        containerProperties.setMessageListener(new Listener<>(consumer));
        containerProperties.setAsyncAcks(false);
        containerProperties.setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
        containerProperties.setGroupId(consumer.group());
        return container;
    }
 
    /**
     * 定義監(jiān)聽
     */
    private static class Listener<T> implements AcknowledgingConsumerAwareMessageListener<String, T> {
        private final KafkaConsumer<T> kafkaConsumer;
 
        public Listener(KafkaConsumer<T> consumer) {
            this.kafkaConsumer = consumer;
        }
 
        @Override
        public void onMessage(ConsumerRecord<String, T> data, Acknowledgment acknowledgment, Consumer<?, ?> consumer) {
            log.info("group【{}】接收到來自topic【{}】的消息", kafkaConsumer.group(), data.topic());
            // 處理數(shù)據(jù)
            kafkaConsumer.process(data.value());
            // 提交offset
            log.info("group【{}】提交topic【{}】的offset", kafkaConsumer.group(), data.topic());
            consumer.commitSync();
        }
    }
}

這個可以根據(jù)需要動態(tài)的啟動消費者

3.2 手動創(chuàng)建KafkaProducer

@Bean
    public KafkaProducer<String, byte[]> kafkaProducer() {
 
        Properties props = new Properties();
 
        // 這里可以配置幾臺broker即可,他會自動從broker去拉取元數(shù)據(jù)進行緩存
        props.put("bootstrap.servers", bootstrapServers);
        // 這個就是負責把發(fā)送的key從字符串序列化為字節(jié)數(shù)組
        props.put("key.serializer", keySerializer);
        // 這個就是負責把你發(fā)送的實際的message從字符串序列化為字節(jié)數(shù)組
        props.put("value.serializer", valueSerializer);
        // 默認是32兆=33554432
        props.put("buffer.memory", bufferMemory);
        // 一般來說是要自己手動設置的,不是純粹依靠默認值的,16kb
        props.put("batch.size", batchSize);
        // 發(fā)送一條消息出去,100ms內(nèi)還沒有湊成一個batch發(fā)送,必須立即發(fā)送出去
        props.put("linger.ms", lingerMs);
        // 這個是說你可以發(fā)送的最大的請求的大小 默認是1m=1048576
//        props.put("max.request.size", 10485760);
        // follower有沒有同步成功你就不管了
        props.put("acks", acks);
        // 這個重試,一般來說,給個3次~5次就足夠了,可以cover住一般的異常場景
        props.put("retries", retries);
        // 每次重試間隔100ms
        props.put("retry.backoff.ms", retryBackOffMs);
 
        props.put("max.in.flight.requests.per.connection", maxInFlightRequestsPerConnection);
 
        return new KafkaProducer<>(props);
    }

4、總結

4.1 區(qū)別

KafkaProducer是Kafka-client提供的原生Java Kafka客戶端發(fā)送消息的API。

KafkaTemplate是Spring Kafka中提供的一個高級工具類,用于可以方便地發(fā)送消息到Kafka。它封裝了KafkaProducer,提供了更多的便利方法和更高級的消息發(fā)送方式。

org.apache.kafka.clients.producer.KafkaProducer

org.springframework.kafka.core.KafkaTemplate

4.2 場景選擇

在spring應用中如果需要訂閱kafka消息,通常情況下我們不會直接使用kafka-client, 而是使用更方便的一層封裝spring-kafka。

不需要動態(tài)的選擇時候可以使用Spring-kafka,在需要動態(tài)創(chuàng)建時可以使用kafka-client的api進行處理

4.3 ConsumerRecord和ProducerRecord

兩者都是kafka-client的類,在Spring-kafka中依然可以使用,可以發(fā)送和接受

以上就是Springboot使用kafka的兩種方式的詳細內(nèi)容,更多關于Springboot使用kafka的資料請關注腳本之家其它相關文章!

相關文章

  • Java 實戰(zhàn)項目錘煉之在線美食網(wǎng)站系統(tǒng)的實現(xiàn)流程

    Java 實戰(zhàn)項目錘煉之在線美食網(wǎng)站系統(tǒng)的實現(xiàn)流程

    讀萬卷書不如行萬里路,只學書上的理論是遠遠不夠的,只有在實戰(zhàn)中才能獲得能力的提升,本篇文章手把手帶你用java+SSM+jsp+mysql+maven實現(xiàn)一個在線美食網(wǎng)站系統(tǒng),大家可以在過程中查缺補漏,提升水平
    2021-11-11
  • spring配置掃描多個包問題解析

    spring配置掃描多個包問題解析

    這篇文章主要介紹了spring配置掃描多個包問題解析,具有一定參考價值,需要的朋友可以了解下。
    2017-10-10
  • javaweb配置jsp路徑映射操作

    javaweb配置jsp路徑映射操作

    這篇文章主要介紹了javaweb配置jsp路徑映射操作,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-08-08
  • Mybatis查詢返回Map<String,Object>類型的實現(xiàn)

    Mybatis查詢返回Map<String,Object>類型的實現(xiàn)

    本文主要介紹了Mybatis查詢返回Map<String,Object>類型的實現(xiàn),文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2023-07-07
  • Java 如何實現(xiàn)AES加密

    Java 如何實現(xiàn)AES加密

    這篇文章主要介紹了Java 如何實現(xiàn)AES加密,幫助大家完成對接,完成自身需求,感興趣的朋友可以了解下
    2020-10-10
  • protobuf與json轉換小結

    protobuf與json轉換小結

    protobuf對象不能直接使用jsonlib去轉,因為protobuf生成的對象的get方法返回的類型有byte[],而只有String類型可以作為json的key,protobuf提供方法進行轉換
    2017-07-07
  • Java實現(xiàn)俄羅斯方塊游戲簡單版

    Java實現(xiàn)俄羅斯方塊游戲簡單版

    這篇文章主要為大家詳細介紹了Java實現(xiàn)俄羅斯方塊游戲簡單版,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-01-01
  • JavaFX中處理Spring的異常的方法

    JavaFX中處理Spring的異常的方法

    JavaFX運行在JavaFX應用線程(JavaFX Application Thread)上,而Spring通常運行在獨立的線程中,本文給大家介紹JavaFX中處理Spring的異常的方法,感興趣的朋友跟隨小編一起看看吧
    2025-10-10
  • java實現(xiàn)科研信息管理系統(tǒng)

    java實現(xiàn)科研信息管理系統(tǒng)

    這篇文章主要為大家詳細介紹了java科研信息管理系統(tǒng),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2018-02-02
  • Spring IOC源碼之bean的注冊過程講解

    Spring IOC源碼之bean的注冊過程講解

    這篇文章主要介紹了Spring IOC源碼之bean的注冊過程講解,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09

最新評論

信丰县| 郴州市| 宝山区| 宁海县| 南靖县| 桃江县| 大足县| 临夏县| 油尖旺区| 资溪县| 通许县| 邢台市| 六枝特区| 鸡泽县| 新乐市| 南靖县| 宁化县| 昌都县| 兰坪| 河池市| 江源县| 清涧县| 微山县| 泰宁县| 得荣县| 榆社县| 五峰| 万宁市| 眉山市| 辽阳县| 石泉县| 同江市| 永和县| 宁远县| 隆昌县| 西乌珠穆沁旗| 武邑县| 临湘市| 隆子县| 西贡区| 哈巴河县|