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

Kafka整合WebFlux實(shí)踐

 更新時(shí)間:2026年03月06日 15:25:28   作者:西紅柿系番茄  
文章介紹了如何在Kafka中整合WebFlux,包括引入依賴(lài)和代碼示例,并對(duì)相關(guān)知識(shí)進(jìn)行了總結(jié)

Kafka整合WebFlux

1、引入依賴(lài)

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<dependency>
    <groupId>io.projectreactor.kafka</groupId>
    <artifactId>reactor-kafka</artifactId>
    <version>1.1.0.RELEASE</version>
</dependency>

2、代碼示例

@Component
public class KafkaService {

    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();

    private KafkaSender<String, String> kafkaSender;
    private KafkaReceiver<String, String> kafkaReceiver;

    @PostConstruct
    public void init() {
        final Map<String, Object> producerProps = new HashMap<>();
        producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
        producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        final SenderOptions<String, String> producerOptions = SenderOptions.create(producerProps);
        this.kafkaSender = KafkaSender.create(producerOptions);

        final Map<String, Object> consumerProps = new HashMap<>();
        consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
        consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        consumerProps.put(ConsumerConfig.CLIENT_ID_CONFIG, "payment-validator-1");
        consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "payment-validator");
        consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        ReceiverOptions<String, String> consumerOptions = ReceiverOptions.<String, String>create(consumerProps)
                .subscription(Collections.singleton("demo"))
                .addAssignListener(partitions -> System.out.println("onPartitionsAssigned " + partitions))
                .addRevokeListener(partitions -> System.out.println("onPartitionsRevoked " + partitions));
        KafkaReceiver<String, String> kafkaReceiver = KafkaReceiver.create(consumerOptions);
        kafkaReceiver.receive().doOnNext(r -> {
            System.out.println(r.value());
            r.receiverOffset().acknowledge();
        }).subscribe();
        this.kafkaReceiver = kafkaReceiver;
    }

    public Mono< ?> send() {
        SenderRecord<String, String, Object> senderRecord = SenderRecord.create(new ProducerRecord<>("demo", value()), 1);
        return kafkaSender.send(Mono.just(senderRecord)).next();
    }

    private String value() {
        Map<String, String> map = new HashMap<>();
        map.put("name", UUID.randomUUID().toString());
        try {
            return OBJECT_MAPPER.writeValueAsString(map);
        } catch (JsonProcessingException e) {
            return "{}";
        }
    }
}

3、其它

server:
  port: 8888

spring:
  jackson:
    serialization:
      FAIL_ON_EMPTY_BEANS: false

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • eclipse報(bào)錯(cuò) eclipse啟動(dòng)報(bào)錯(cuò)解決方法

    eclipse報(bào)錯(cuò) eclipse啟動(dòng)報(bào)錯(cuò)解決方法

    本文將介紹eclipse啟動(dòng)報(bào)錯(cuò)解決方法,需要了解的朋友可以參考下
    2012-11-11
  • java8 實(shí)現(xiàn)提取集合對(duì)象的每個(gè)屬性

    java8 實(shí)現(xiàn)提取集合對(duì)象的每個(gè)屬性

    這篇文章主要介紹了java8 實(shí)現(xiàn)提取集合對(duì)象的每個(gè)屬性方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2021-02-02
  • Golang Protocol Buffer案例詳解

    Golang Protocol Buffer案例詳解

    這篇文章主要介紹了Golang Protocol Buffer案例詳解,本篇文章通過(guò)簡(jiǎn)要的案例,講解了該項(xiàng)技術(shù)的了解與使用,以下就是詳細(xì)內(nèi)容,需要的朋友可以參考下
    2021-08-08
  • Java中的StopWatch計(jì)時(shí)利器使用指南

    Java中的StopWatch計(jì)時(shí)利器使用指南

    StopWatch通常用于測(cè)量一段代碼執(zhí)行所花費(fèi)的時(shí)間,它能夠精確地記錄開(kāi)始時(shí)間、結(jié)束時(shí)間,并計(jì)算出這中間的時(shí)間差,下面給大家介紹Java中的StopWatch計(jì)時(shí)利器的深度解析與使用指南,感興趣的朋友一起看看吧
    2025-05-05
  • Spring?Boot中記錄用戶(hù)系統(tǒng)操作流程

    Spring?Boot中記錄用戶(hù)系統(tǒng)操作流程

    這篇文章主要介紹了如何在Spring?Boot中記錄用戶(hù)系統(tǒng)操作流程,將介紹如何在Spring?Boot中使用AOP(面向切面編程)和日志框架來(lái)實(shí)現(xiàn)用戶(hù)系統(tǒng)操作流程的記錄,需要的朋友可以參考下
    2023-07-07
  • idea中Tomcat啟動(dòng)失敗的解決

    idea中Tomcat啟動(dòng)失敗的解決

    這篇文章主要介紹了idea中Tomcat啟動(dòng)失敗的解決,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-09-09
  • Mybatis 動(dòng)態(tài)SQL的幾種實(shí)現(xiàn)方法

    Mybatis 動(dòng)態(tài)SQL的幾種實(shí)現(xiàn)方法

    這篇文章主要介紹了Mybatis 動(dòng)態(tài)SQL的幾種實(shí)現(xiàn)方法,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2020-11-11
  • Java Kafka實(shí)現(xiàn)延遲隊(duì)列的示例代碼

    Java Kafka實(shí)現(xiàn)延遲隊(duì)列的示例代碼

    kafka作為一個(gè)使用廣泛的消息隊(duì)列,很多人都不會(huì)陌生。本文將利用Kafka實(shí)現(xiàn)延遲隊(duì)列,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以嘗試一下
    2022-08-08
  • Java圖論的兩個(gè)基本概念之有向圖與無(wú)向圖詳解

    Java圖論的兩個(gè)基本概念之有向圖與無(wú)向圖詳解

    圖論是數(shù)學(xué)的一個(gè)基本分支,涉及對(duì)圖研究,圖是復(fù)雜數(shù)據(jù)結(jié)構(gòu)的可視化表示,有助于理解不同實(shí)體之間的關(guān)系,這篇文章主要介紹了Java圖論的兩個(gè)基本概念之有向圖與無(wú)向圖的相關(guān)資料,需要的朋友可以參考下
    2026-03-03
  • java使用Jco連接SAP過(guò)程

    java使用Jco連接SAP過(guò)程

    這篇文章主要介紹了java使用Jco連接SAP過(guò)程,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-03-03

最新評(píng)論

定州市| 宁晋县| 武强县| 辽中县| 镇康县| 盐边县| 榆林市| 江城| 茂名市| 毕节市| 揭阳市| 新蔡县| 莲花县| 资溪县| 鹤庆县| 扎囊县| 云阳县| 磴口县| 文山县| 将乐县| 呼伦贝尔市| 宽甸| 余庆县| 吴川市| 阿荣旗| 永嘉县| 固始县| 手游| 乐昌市| 乐至县| 筠连县| 灌阳县| 镇沅| 富阳市| 陈巴尔虎旗| 册亨县| 石棉县| 凤凰县| 定州市| 大荔县| 滕州市|