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

springboot使用@KafkaListener監(jiān)聽多個kafka配置實現

 更新時間:2024年04月09日 09:36:45   作者:道不平  
當服務中需要監(jiān)聽多個kafka時,?需要配置多個kafka,本文主要介紹了springboot使用@KafkaListener監(jiān)聽多個kafka配置實現,具有一定的參考價值,感興趣的可以了解一下

背景

使用springboot整合kafka時, springboot默認讀取配置文件中 spring.kafka...配置初始化kafka, 使用@KafkaListener時指定topic即可, 當服務中需要監(jiān)聽多個kafka時, 需要配置多個kafka, 這種方式不適用

方案

可以手動讀取不同kafka配置信息, 創(chuàng)建不同的Kafka 監(jiān)聽容器工廠, 使用@KafkaListener時指定相應的容器工廠, 代碼如下:

1. 導入依賴

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

2. yml配置

kafka:
  # 默認消費者配置
  default-consumer:
    # 自動提交已消費offset
    enable-auto-commit: true
    # 自動提交間隔時間
    auto-commit-interval: 1000
    # 消費的超時時間
    poll-timeout: 1500
    # 如果Kafka中沒有初始偏移量,或者服務器上不再存在當前偏移量(例如,因為該數據已被刪除)自動將該偏移量重置成最新偏移量
    auto.offset.reset: latest
    # 消費會話超時時間(超過這個時間consumer沒有發(fā)送心跳,就會觸發(fā)rebalance操作)
    session.timeout.ms: 120000
    # 消費請求超時時間
    request.timeout.ms: 180000
  # 1號kafka配置
  test1:
    bootstrap-servers: xxxx:xxxx,xxxx:xxxx,xxxx:xxxx
    consumer:
      group-id: xxx
      sasl.mechanism: xxxx
      security.protocol: xxxx
      sasl.jaas.config: xxxx
  # 2號kafka配置
  test2:
    bootstrap-servers: xxxx:xxxx,xxxx:xxxx,xxxx:xxxx
    consumer:
      group-id: xxx
      sasl.mechanism: xxxx
      security.protocol: xxxx
      sasl.jaas.config: xxxx

3. 容器工廠配置

package com.zhdx.modules.backstage.config;

import com.google.common.collect.Maps;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.ConcurrentMessageListenerContainer;

import java.util.Map;

/**
 * kafka監(jiān)聽容器工廠配置
 * <p>
 * 拓展其他消費者配置只需配置指定的屬性和bean即可
 */
@EnableKafka
@Configuration
@RefreshScope
public class KafkaListenerContainerFactoryConfig {

    /**
     *  test1 kafka配置
     */
    @Value("${kafka.test1.bootstrap-servers}")
    private String test1KafkaServerUrls;

    @Value("${kafka.test1.consumer.group-id}")
    private String test1GroupId;

    @Value("${kafka.test1.consumer.sasl.mechanism}")
    private String test1SaslMechanism;

    @Value("${kafka.test1.consumer.security.protocol}")
    private String test1SecurityProtocol;

    @Value("${kafka.test1.consumer.sasl.jaas.config}")
    private String test1SaslJaasConfig;
    /**
     *  test2 kafka配置
     */
    @Value("${kafka.test2.bootstrap-servers}")
    private String test2KafkaServerUrls;

    @Value("${kafka.test2.consumer.group-id}")
    private String test2GroupId;

    @Value("${kafka.test2.consumer.sasl.mechanism}")
    private String test2SaslMechanism;

    @Value("${kafka.test2.consumer.security.protocol}")
    private String test2SecurityProtocol;

    @Value("${kafka.test2.consumer.sasl.jaas.config}")
    private String test2SaslJaasConfig;

    /**
     * 默認消費者配置
     */
    @Value("${kafka.default-consumer.enable-auto-commit}")
    private boolean enableAutoCommit;

    @Value("${kafka.default-consumer.poll-timeout}")
    private int pollTimeout;

    @Value("${kafka.default-consumer.auto.offset.reset}")
    private String autoOffsetReset;

    @Value("${kafka.default-consumer.session.timeout.ms}")
    private int sessionTimeoutMs;

    @Value("${kafka.default-consumer.request.timeout.ms}")
    private int requestTimeoutMs;

    /**
     * test1消費者配置
     */
    public Map<String, Object> test1ConsumerConfigs() {
        Map<String, Object> props = getDefaultConsumerConfigs();
        // broker server地址
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, test1KafkaServerUrls);
        // 消費者組
        props.put(ConsumerConfig.GROUP_ID_CONFIG, test1GroupId);
        // 加密
        props.put(SaslConfigs.SASL_MECHANISM, test1SaslMechanism);
        props.put("security.protocol", test1SecurityProtocol);
        // 賬號密碼
        props.put(SaslConfigs.SASL_JAAS_CONFIG, test1SaslJaasConfig);
        return props;
    }
    
    /**
     * test2消費者配置
     */
    public Map<String, Object> test2ConsumerConfigs() {
        Map<String, Object> props = getDefaultConsumerConfigs();
        // broker server地址
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, test2KafkaServerUrls);
        // 消費者組
        props.put(ConsumerConfig.GROUP_ID_CONFIG, test2GroupId);
        // 加密
        props.put(SaslConfigs.SASL_MECHANISM, test2SaslMechanism);
        props.put("security.protocol", test2SecurityProtocol);
        // 賬號密碼
        props.put(SaslConfigs.SASL_JAAS_CONFIG, test2SaslJaasConfig);
        return props;
    }

    /**
     * 默認消費者配置
     */
    private Map<String, Object> getDefaultConsumerConfigs() {
        Map<String, Object> props = Maps.newHashMap();
        // 自動提交(按周期)已消費offset 批量消費下設置false
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, enableAutoCommit);
        // 消費會話超時時間(超過這個時間consumer沒有發(fā)送心跳,就會觸發(fā)rebalance操作)
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, sessionTimeoutMs);
        // 消費請求超時時間
        props.put(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG, requestTimeoutMs);
        // 序列化
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // 如果Kafka中沒有初始偏移量,或者服務器上不再存在當前偏移量(例如,因為該數據已被刪除)自動將該偏移量重置成最新偏移量
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset);
        return props;
    }

    /**
     * 消費者工廠類
     */
    public ConsumerFactory<String, String> initConsumerFactory(Map<String, Object> consumerConfigs) {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs);
    }

    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> initKafkaListenerContainerFactory(
        Map<String, Object> consumerConfigs) {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(initConsumerFactory(consumerConfigs));
        // 是否開啟批量消費
        factory.setBatchListener(false);
        // 消費的超時時間
        factory.getContainerProperties().setPollTimeout(pollTimeout);
        return factory;
    }

    /**
     * 創(chuàng)建test1 Kafka 監(jiān)聽容器工廠。
     *
     * @return KafkaListenerContainerFactory<ConcurrentMessageListenerContainer < String, String>> 返回的 KafkaListenerContainerFactory 對象
     */
    @Bean(name = "test1KafkaListenerContainerFactory")
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> test1KafkaListenerContainerFactory() {
        Map<String, Object> consumerConfigs = this.test1ConsumerConfigs();
        return initKafkaListenerContainerFactory(consumerConfigs);
    }
    

    /**
     * 創(chuàng)建test2 Kafka 監(jiān)聽容器工廠。
     *
     * @return KafkaListenerContainerFactory<ConcurrentMessageListenerContainer < String, String>> 返回的 KafkaListenerContainerFactory 對象
     */
    @Bean(name = "test2KafkaListenerContainerFactory")
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> test2KafkaListenerContainerFactory() {
        Map<String, Object> consumerConfigs = this.test2ConsumerConfigs();
        return initKafkaListenerContainerFactory(consumerConfigs);
    }
}

4. @KafkaListener使用

package com.zhdx.modules.backstage.kafka;

import com.alibaba.fastjson.JSON;

import lombok.extern.slf4j.Slf4j;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

/**
 * kafka監(jiān)聽器
 */
@Slf4j
@Component
public class test1KafkaListener {
    @KafkaListener(containerFactory = "test1KafkaListenerContainerFactory", topics = "xxx")
    public void handleHyPm(ConsumerRecord<String, String> record) {
        log.info("消費到topic xxx消息:{}", JSON.toJSONString(record.value()));
    }
}

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

相關文章

  • Springboot升級到2.7.2結合nacos遇到的坑及解決

    Springboot升級到2.7.2結合nacos遇到的坑及解決

    這篇文章主要介紹了Springboot升級到2.7.2結合nacos遇到的坑及解決,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-06-06
  • 使用JDBC4.0操作XML類型的字段(保存獲取xml數據)的方法

    使用JDBC4.0操作XML類型的字段(保存獲取xml數據)的方法

    jdbc4.0最重要的特征是支持xml數據類型,接下來通過本文重點給大家介紹如何使用jdbc4.0操作xml類型的字段,對jdbc4.0 xml相關知識感興趣的朋友一起看下吧
    2016-08-08
  • SpringBoot切面攔截@PathVariable參數及拋出異常的全局處理方式

    SpringBoot切面攔截@PathVariable參數及拋出異常的全局處理方式

    這篇文章主要介紹了SpringBoot切面攔截@PathVariable參數及拋出異常的全局處理方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • java.lang.Void 與 void的比較及使用方法介紹

    java.lang.Void 與 void的比較及使用方法介紹

    這篇文章主要介紹了java.lang.Void 與 void的比較及使用方法介紹,小編覺得挺不錯的,這里給大家分享一下,需要的朋友可以參考。
    2017-10-10
  • Java HttpURLConnection使用方法詳解

    Java HttpURLConnection使用方法詳解

    這篇文章主要為大家詳細介紹了Java HttpURLConnection使用方法,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-11-11
  • Java使用System.currentTimeMillis()方法計算程序運行時間的示例代碼

    Java使用System.currentTimeMillis()方法計算程序運行時間的示例代碼

    System.currentTimeMillis() 方法的返回類型為 long ,表示毫秒為單位的當前時間,文中通過示例代碼介紹了計算 String 類型與 StringBuilder 類型拼接字符串的耗時情況,對Java計算程序運行時間相關知識感興趣的朋友一起看看吧
    2022-03-03
  • 使用Spring Cache設置緩存條件操作

    使用Spring Cache設置緩存條件操作

    這篇文章主要介紹了使用Spring Cache設置緩存條件操作,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • Java并發(fā)Futures和Callables類實例詳解

    Java并發(fā)Futures和Callables類實例詳解

    Callable對象返回Future對象,該對象提供監(jiān)視線程執(zhí)行的任務進度的方法, Future對象可用于檢查Callable的狀態(tài),然后線程完成后從Callable中檢索結果,這篇文章給大家介紹Java并發(fā)Futures和Callables類的相關知識,感興趣的朋友一起看看吧
    2024-05-05
  • java去除空格、標點符號的方法實例

    java去除空格、標點符號的方法實例

    這篇文章主要給大家介紹了關于java去除空格、標點符號的相關資料,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2020-09-09
  • Java中實現多線程關鍵詞整理(總結)

    Java中實現多線程關鍵詞整理(總結)

    這篇文章主要介紹了Java中實現多線程關鍵詞整理,非常不錯,具有參考借鑒價值,需要的朋友可以參考下
    2017-05-05

最新評論

宜丰县| 南陵县| 枣阳市| 麻栗坡县| 仁化县| 安顺市| 天全县| 北海市| 定南县| 炉霍县| 彭山县| 鹤峰县| 兰西县| 清河县| 金堂县| 潢川县| 新绛县| 阿克陶县| 天水市| 盐津县| 丹阳市| 灵寿县| 监利县| 石泉县| 牡丹江市| 盈江县| 濮阳市| 陇西县| 蓬莱市| 奈曼旗| 招远市| 绵阳市| 中超| 鄂托克前旗| 定南县| 白玉县| 高邑县| 武邑县| 邳州市| 平邑县| 攀枝花市|