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

springboot連接kafka集群的使用示例

 更新時(shí)間:2023年09月14日 10:57:35   作者:timi先生  
在項(xiàng)目中使用kafka的場(chǎng)景有很多,尤其是實(shí)時(shí)產(chǎn)生的數(shù)據(jù)流,本文主要介紹了springboot連接kafka集群的使用示例,具有一定的參考價(jià)值,感興趣的可以了解一下

一、環(huán)境搭建

1.1 springboot 環(huán)境

  • JDK 11+
  • Maven 3.8.x+
  • springboot 2.5.4 +

1.2 kafka 依賴

springboot的pom文件導(dǎo)入

       <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka-test</artifactId>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.kafka</groupId>
            <artifactId>kafka-clients</artifactId>
            <version>3.4.0</version>
        </dependency>

二、 kafka 配置類

2.1 發(fā)布者

2.1.1 配置

發(fā)布者我們使用 KafkaTemplate 來(lái)進(jìn)行消息發(fā)布,所以需要先對(duì)其進(jìn)行一些必要的配置。

@Configuration
@EnableKafka
public class KafkaConfig {
     /***** 發(fā)布者 *****/
    //生產(chǎn)者工廠
    @Bean
    public ProducerFactory<Integer, String> producerFactory() {
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }
    //生產(chǎn)者配置
    @Bean
    public Map<String, Object> producerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.2.83:9092,192.168.2.84:9092,192.168.2.86:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return props;
    }
    //生產(chǎn)者模板
    @Bean
    public KafkaTemplate<Integer, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

2.1.2 構(gòu)建發(fā)布者類

配置完發(fā)布者,下來(lái)就是發(fā)布消息,我們需要繼承 ProducerListener<K, V> 接口,該接口完整信息如下:

public interface ProducerListener<K, V> {
    void onSuccess(ProducerRecord<K, V> producerRecord, RecordMetadata recordMetadata);
    void onError(ProducerRecord<K, V> producerRecord, RecordMetadata recordMetadata,
            Exception exception);
}

實(shí)現(xiàn)該接口的方法,我們可以獲取包含發(fā)送結(jié)果(成功或失?。┑漠惒交卣{(diào),也就是可以在這個(gè)接口的實(shí)現(xiàn)中獲取發(fā)送結(jié)果。

我們簡(jiǎn)單的實(shí)現(xiàn)構(gòu)建一個(gè)發(fā)布者類,接收主題和發(fā)布消息參數(shù),并打印發(fā)布結(jié)果。

@Component
public class KafkaProducer implements ProducerListener<Object,Object> {
    private static final Logger producerlog = LoggerFactory.getLogger(KafkaProducer.class);
    private final KafkaTemplate<Integer, String> kafkaTemplate;
    public KafkaProducer(KafkaTemplate<Integer, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }
    public void producer (String msg,String topic){
        ListenableFuture<SendResult<Integer, String>> future = kafkaTemplate.send(topic,0, msg);
        future.addCallback(new KafkaSendCallback<Integer, String>() {
            @Override
            public void onSuccess(SendResult<Integer, String> result) {
                producerlog.info("發(fā)送成功 {}", result);
            }
            @Override
            public void onFailure(KafkaProducerException ex) {
                ProducerRecord<Integer, String> failed = ex.getFailedProducerRecord();
                producerlog.info("發(fā)送失敗 {}",failed);
            }
        });
    }
}

2.1.3 發(fā)布消息

寫一個(gè)controller類來(lái)測(cè)試我們構(gòu)建的發(fā)布者類,這個(gè)類中打印接收到的消息,來(lái)確保信息接收不出問(wèn)題。

@RestController
public class KafkaTestController {
    private static final Logger kafkaTestLog = LoggerFactory.getLogger(KafkaTestController.class);
    @Resource
    private KafkaProducer kafkaProducer;
    @GetMapping("/kafkaTest")
    public void kafkaTest(String msg,String topic){
        kafkaProducer.producer(msg,topic);
        kafkaTestLog.info("接收到消息 {} {}",msg,topic);
    }
}

一切準(zhǔn)備就緒,我們啟動(dòng)程序利用postman來(lái)進(jìn)行簡(jiǎn)單的測(cè)試。

進(jìn)行消息發(fā)布:

發(fā)布結(jié)果:

可以看到消息發(fā)送成功。

我們?cè)倏纯磌afka消費(fèi)者有沒(méi)有接收到消息:

看以看到,kakfa的消費(fèi)者也接收到了消息。

2.2 消費(fèi)者

2.2.1 配置

消息的接受有多種方式,我們這里選擇的是使用 @KafkaListener 注解來(lái)進(jìn)行消息接收。它的使用像下面這樣:

public class Listener {
    @KafkaListener(id = "foo", topics = "myTopic", clientIdPrefix = "myClientId")
    public void listen(String data) {
        ...
    }
}

看起來(lái)不是太難吧,但使用這個(gè)注解,我們需要配置底層 ConcurrentMessageListenerContainer.kafkaListenerContainerFactor。

我們?cè)谠瓉?lái)的kafka配置類 KafkaConfig 中,繼續(xù)配置消費(fèi)者,大概就像下面這樣

@Configuration
@EnableKafka
public class KafkaConfig {
     /***** 發(fā)布者 *****/
    //生產(chǎn)者工廠
    @Bean
    public ProducerFactory<Integer, String> producerFactory() {
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }
    //生產(chǎn)者配置
    @Bean
    public Map<String, Object> producerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.2.83:9092,192.168.2.84:9092,192.168.2.86:9092");
           props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return props;
    }
    //生產(chǎn)者模板
    @Bean
    public KafkaTemplate<Integer, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
    /***** 消費(fèi)者 *****/
    //容器監(jiān)聽工廠
    @Bean
    KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>>
    kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<Integer, String> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setConcurrency(3);
        factory.getContainerProperties().setPollTimeout(3000);
        return factory;
    }
    //消費(fèi)者工廠
    @Bean
    public ConsumerFactory<Integer, String> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }
    //消費(fèi)者配置
    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,"192.168.2.83:9092,192.168.2.84:9092,192.168.2.86:9092");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, JsonDeserializer.class);
        props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName());
        props.put(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG,3000);
        return props;
    }
}

注意,要設(shè)置容器屬性必須使用getContainerProperties()工廠方法。它用作注入容器的實(shí)際屬性的模板

2.2.2 構(gòu)建消費(fèi)者類

配置好后,我們就可以使用這個(gè)注解了。這個(gè)注解的使用有多種方式:

1、用它來(lái)覆蓋容器工廠的concurrency和屬性

@KafkaListener(id = "myListener", topics = "myTopic",
        autoStartup = "${listen.auto.start:true}", concurrency = "${listen.concurrency:3}")
public void listen(String data) {
    ...
}

2、可以使用顯式主題和分區(qū)(以及可選的初始偏移量)

@KafkaListener(id = "thing2", topicPartitions =
        { @TopicPartition(topic = "topic1", partitions = { "0", "1" }),
          @TopicPartition(topic = "topic2", partitions = "0",
             partitionOffsets = @PartitionOffset(partition = "1", initialOffset = "100"))
        })
public void listen(ConsumerRecord<?, ?> record) {
    ...
}

3、將初始偏移應(yīng)用于所有已分配的分區(qū)

@KafkaListener(id = "thing3", topicPartitions =
        { @TopicPartition(topic = "topic1", partitions = { "0", "1" },
             partitionOffsets = @PartitionOffset(partition = "*", initialOffset = "0"))
        })
public void listen(ConsumerRecord<?, ?> record) {
    ...
}

4、指定以逗號(hào)分隔的分區(qū)列表或分區(qū)范圍

@KafkaListener(id = "pp", autoStartup = "false",
        topicPartitions = @TopicPartition(topic = "topic1",
                partitions = "0-5, 7, 10-15"))
public void process(String in) {
    ...
}

5、可以向偵聽器提供Acknowledgment

@KafkaListener(id = "cat", topics = "myTopic",
          containerFactory = "kafkaManualAckListenerContainerFactory")
public void listen(String data, Acknowledgment ack) {
    ...
    ack.acknowledge();
}

6、添加標(biāo)頭

@KafkaListener(id = "list", topics = "myTopic", containerFactory = "batchFactory")
public void listen(List<String> list,
        @Header(KafkaHeaders.RECEIVED_KEY) List<Integer> keys,
        @Header(KafkaHeaders.RECEIVED_PARTITION) List<Integer> partitions,
        @Header(KafkaHeaders.RECEIVED_TOPIC) List<String> topics,
        @Header(KafkaHeaders.OFFSET) List<Long> offsets) {
    ...
}

我們這里寫一個(gè)簡(jiǎn)單的,只用它來(lái)接受指定主題的數(shù)據(jù):

@Component
public class KafkaConsumer {
    private static final Logger consumerlog = LoggerFactory.getLogger(KafkaConsumer.class);
    @KafkaListener(topicPartitions  = @TopicPartition(topic = "kafka-topic-test",
            partitions = "0"))
    public void consumer (String data){
        consumerlog.info("消費(fèi)者接收數(shù)據(jù) {}",data);
    }
}

這里解釋一下,因?yàn)槲覀冞M(jìn)行了手動(dòng)分配主題/分區(qū),所以 注解中g(shù)roup.id 可以為空。若要指定group.id請(qǐng)?jiān)谙M(fèi)者配置中加上props.put(ConsumerConfig.GROUP_ID_CONFIG, “bzt001”); 或在 @TopicPartition 注解后加上 groupId = “組id”

2.2.3 進(jìn)行消息消費(fèi)

繼續(xù)使用postman調(diào)用我們寫好的發(fā)布者發(fā)布消息,觀察控制臺(tái)的消費(fèi)者類是否有相關(guān)日志出現(xiàn)。

 到此這篇關(guān)于springboot連接kafka集群的使用示例的文章就介紹到這了,更多相關(guān)springboot連接kafka集群內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java中for與foreach的區(qū)別

    Java中for與foreach的區(qū)別

    本文主要介紹了Java中for與foreach的區(qū)別,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2022-05-05
  • springboot中如何通過(guò)cors協(xié)議解決跨域問(wèn)題

    springboot中如何通過(guò)cors協(xié)議解決跨域問(wèn)題

    這篇文章主要介紹了springboot中通過(guò)cors協(xié)議解決跨域問(wèn)題,cors是一個(gè)w3c標(biāo)準(zhǔn),它允許瀏覽器(目前ie8以下還不能被支持)像我們不同源的服務(wù)器發(fā)出xmlHttpRequest請(qǐng)求,我們可以繼續(xù)使用ajax進(jìn)行請(qǐng)求訪問(wèn)。具體內(nèi)容詳情大家跟隨腳本之家小編一起學(xué)習(xí)吧
    2018-05-05
  • 使用Jenv管理多版本JDK環(huán)境的詳細(xì)教程

    使用Jenv管理多版本JDK環(huán)境的詳細(xì)教程

    在現(xiàn)代 Java 開發(fā)中,我們經(jīng)常需要在不同的項(xiàng)目中使用不同版本的 JDK,手動(dòng)切換 JAVA_HOME 環(huán)境變量既繁瑣又容易出錯(cuò),Jenv 是一個(gè)優(yōu)秀的 JDK 版本管理工具,可以讓我們輕松地在不同 JDK 版本間切換,所以本文給大家介紹了使用Jenv管理多版本JDK環(huán)境的詳細(xì)教程
    2025-08-08
  • Java使用File類遍歷目錄及文件實(shí)例代碼

    Java使用File類遍歷目錄及文件實(shí)例代碼

    本篇文章主要介紹了Java使用File類遍歷目錄及文件實(shí)例代碼,詳細(xì)的介紹了File類的使用,有興趣的可以了解一下。
    2017-04-04
  • Struts2在打包json格式的懶加載異常問(wèn)題

    Struts2在打包json格式的懶加載異常問(wèn)題

    這篇文章主要為大家詳細(xì)介紹了Struts2在打包json格式的懶加載異常問(wèn)題,感興趣的小伙伴們可以參考一下
    2016-06-06
  • Java實(shí)現(xiàn)并查集

    Java實(shí)現(xiàn)并查集

    這篇文章主要為大家詳細(xì)介紹了Java實(shí)現(xiàn)并查集,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2020-07-07
  • java前后端使用ajax數(shù)據(jù)交互問(wèn)題(簡(jiǎn)單demo)

    java前后端使用ajax數(shù)據(jù)交互問(wèn)題(簡(jiǎn)單demo)

    這篇文章主要介紹了java前后端使用ajax數(shù)據(jù)交互問(wèn)題(簡(jiǎn)單demo),具有很好的參考價(jià)值,希望對(duì)大家有所幫助。
    2023-06-06
  • 詳解Java分布式緩存系統(tǒng)中必須解決的四大問(wèn)題

    詳解Java分布式緩存系統(tǒng)中必須解決的四大問(wèn)題

    分布式緩存系統(tǒng)是三高架構(gòu)中不可或缺的部分,極大地提高了整個(gè)項(xiàng)目的并發(fā)量、響應(yīng)速度,但它也帶來(lái)了新的需要解決的問(wèn)題,分別是: 緩存穿透、緩存擊穿、緩存雪崩和緩存一致性問(wèn)題。本文將詳細(xì)講解一下這四大問(wèn)題,需要的可以參考一下
    2022-04-04
  • Spring依賴注入的兩種方式(根據(jù)實(shí)例詳解)

    Spring依賴注入的兩種方式(根據(jù)實(shí)例詳解)

    這篇文章主要介紹了Spring依賴注入的兩種方式(根據(jù)實(shí)例詳解),非常具有實(shí)用價(jià)值,需要的朋友可以參考下
    2017-05-05
  • Java中輸入與輸出的方法總結(jié)

    Java中輸入與輸出的方法總結(jié)

    這篇文章主要為大家總結(jié)了一下Java中輸入與輸出的三種方法,文中通過(guò)示例詳細(xì)的講解了一下這些方法的使用,需要的小伙伴可以參考一下
    2022-04-04

最新評(píng)論

策勒县| 鹤峰县| 枣强县| 永和县| 新建县| 台安县| 时尚| 茌平县| 垣曲县| 鲁山县| 古浪县| 贡觉县| 渑池县| 巴东县| 凌海市| 石景山区| 北流市| 灵寿县| 军事| 新闻| 叙永县| 甘洛县| 安宁市| 济阳县| 米泉市| 长寿区| 额济纳旗| 深水埗区| 青神县| 天峻县| 淳安县| 安塞县| 江华| 建瓯市| 滨州市| 微山县| 洞头县| 浦东新区| 穆棱市| 化德县| 海安县|