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

Spring?Boot整合Kafka+SSE實(shí)現(xiàn)實(shí)時(shí)數(shù)據(jù)展示

 更新時(shí)間:2024年06月14日 08:28:10   作者:艾迪的技術(shù)之路  
本文主要介紹了Spring?Boot整合Kafka+SSE實(shí)現(xiàn)實(shí)時(shí)數(shù)據(jù)展示

為什么使用Kafka?

不使用Rabbitmq或者Rocketmq是因?yàn)镵afka是Hadoop集群下的組成部分,對于大數(shù)據(jù)的相關(guān)開發(fā)適應(yīng)性好,且當(dāng)前業(yè)務(wù)場景下不需要使用死信隊(duì)列,不過要注意Kafka對于更新時(shí)間慢的數(shù)據(jù)拉取也較慢,因此對與實(shí)時(shí)性要求高可以選擇其他MQ。

使用消息隊(duì)列是因?yàn)樵撝虚g件具有實(shí)時(shí)性,且可以作為廣播進(jìn)行消息分發(fā)。

為什么使用SSE?

使用Websocket傳輸信息的時(shí)候,會(huì)轉(zhuǎn)成二進(jìn)制數(shù)據(jù),產(chǎn)生一定的時(shí)間損耗,SSE直接傳輸文本,不存在這個(gè)問題

由于Websocket是雙向的,讀取日志的時(shí)候,如果有人連接ws日志,會(huì)發(fā)送大量異常信息,會(huì)給使用段和日志段造成問題;SSE是單向的,不需要考慮這個(gè)問題,提高了安全性
另外就是SSE支持?jǐn)嗑€重連;Websocket協(xié)議本身并沒有提供心跳機(jī)制,所以長時(shí)間沒有數(shù)據(jù)發(fā)送時(shí),會(huì)將這個(gè)連接斷掉,因此需要手寫心跳機(jī)制進(jìn)行實(shí)現(xiàn)。

此外,由于是長連接的一個(gè)實(shí)現(xiàn)方式,所以SSE也可以替代Websocket實(shí)現(xiàn)掃碼登陸(比如通過SSE的超時(shí)組件在實(shí)現(xiàn)二維碼的超時(shí)功能,具體實(shí)現(xiàn)我可以整理一下)

另外,如果是普通項(xiàng)目,不需要過高的實(shí)時(shí)性,則不需要使用Websocket,使用SSE即可

代碼實(shí)現(xiàn)

pom.xml引入SSE和Kafka

<!-- SSE,一般springboot開發(fā)web應(yīng)用的都有 -->
       <dependency>
           <groupId>org.springframework.boot</groupId>
           <artifactId>spring-boot-starter-web</artifactId>
       </dependency>
<!-- kafka,最主要的是第一個(gè),剩下兩個(gè)是測試用的 -->
       <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>

application.properties增加Kafka配置信息

# KafkaProperties
spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.consumer.group-id=community-consumer-group
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer

配置Kafka信息

@Configuration
public class KafkaProducerConfig {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    public Map<String, Object> producerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return props;
    }

    @Bean
    public ProducerFactory<String, String> producerFactory() {
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }

}

配置controller,通過web方式開啟效果

@RestController
@RequestMapping(path = "sse")
public class KafkaSSEController {

    private static final Map<String, SseEmitter> sseCache = new ConcurrentHashMap<>();

    @Resource
    private KafkaTemplate<String, String> kafkaTemplate;

    /**
     * @param message
     * @apiNote 發(fā)送信息到Kafka主題中
     */
    @PostMapping("/send")
    public void sendMessage(@RequestBody String message) {
        kafkaTemplate.send("my-topic", message);
    }

    /**
     * 監(jiān)聽Kafka數(shù)據(jù)
     *
     * @param message
     */
    @KafkaListener(topics = "my-topic", groupId = "my-group-id")
    public void consume(String message) {
        System.out.println("Received message: " + message);
    }

    /**
     * 連接sse服務(wù)
     *
     * @param id
     * @return
     * @throws IOException
     */
    @GetMapping(path = "subscribe", produces = {MediaType.TEXT_EVENT_STREAM_VALUE})
    public SseEmitter push(@RequestParam("id") String id) throws IOException {
        // 超時(shí)時(shí)間設(shè)置為5分鐘,用于演示客戶端自動(dòng)重連
        SseEmitter sseEmitter = new SseEmitter(5_60_000L);
        // 設(shè)置前端的重試時(shí)間為1s
        // send(): 發(fā)送數(shù)據(jù),如果傳入的是一個(gè)非SseEventBuilder對象,那么傳遞參數(shù)會(huì)被封裝到 data 中
        sseEmitter.send(SseEmitter.event().reconnectTime(1000).data("連接成功"));
        sseCache.put(id, sseEmitter);
        System.out.println("add " + id);
        sseEmitter.send("你好", MediaType.APPLICATION_JSON);
        SseEmitter.SseEventBuilder data = SseEmitter.event().name("finish").id("6666").data("哈哈");
        sseEmitter.send(data);
        // onTimeout(): 超時(shí)回調(diào)觸發(fā)
        sseEmitter.onTimeout(() -> {
            System.out.println(id + "超時(shí)");
            sseCache.remove(id);
        });
        // onCompletion(): 結(jié)束之后的回調(diào)觸發(fā)
        sseEmitter.onCompletion(() -> System.out.println("完成?。?!"));
        return sseEmitter;
    }
    /**
     * http://127.0.0.1:8080/sse/push?id=7777&content=%E4%BD%A0%E5%93%88aaaaaa
     * @param id
     * @param content
     * @return
     * @throws IOException
     */
    @ResponseBody
    @GetMapping(path = "push")
    public String push(String id, String content) throws IOException {
        SseEmitter sseEmitter = sseCache.get(id);
        if (sseEmitter != null) {
            sseEmitter.send(content);
        }
        return "over";
    }

    @ResponseBody
    @GetMapping(path = "over")
    public String over(String id) {
        SseEmitter sseEmitter = sseCache.get(id);
        if (sseEmitter != null) {
            // complete(): 表示執(zhí)行完畢,會(huì)斷開連接
            sseEmitter.complete();
            sseCache.remove(id);
        }
        return "over";
    }

}

前端方式

<html>
  <head>
    <script>
      console.log('start')
      const clientId = "your_client_id_x"; // 設(shè)置客戶端ID
      const eventSource = new EventSource(`http://localhost:9999/v1/sse/subscribe/${clientId}`); // 訂閱服務(wù)器端的SSE

      eventSource.onmessage = event => {
        console.log(event.data)
        const message = JSON.parse(event.data);
        console.log(`Received message from server: ${message}`);
      };

      // 發(fā)送消息給服務(wù)器端 可通過 postman 調(diào)用,所以下面 sendMessage() 調(diào)用被注釋掉了
      function sendMessage() {
        const message = "hello sse";
        fetch(`http://localhost:9999/v1/sse/publish/${clientId}`, {
          method: "POST",
          headers: { "Content-Type": "application/json" },
          body: JSON.stringify(message)
        });
        console.log('dddd'+JSON.stringify(message))
      }
      // sendMessage()
    </script>
  </head>
</html>

到此這篇關(guān)于Spring Boot整合Kafka+SSE實(shí)現(xiàn)實(shí)時(shí)數(shù)據(jù)展示的文章就介紹到這了,更多相關(guān)SpringBoot實(shí)時(shí)數(shù)據(jù)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java學(xué)習(xí)筆記:基本輸入、輸出數(shù)據(jù)操作實(shí)例分析

    Java學(xué)習(xí)筆記:基本輸入、輸出數(shù)據(jù)操作實(shí)例分析

    這篇文章主要介紹了Java學(xué)習(xí)筆記:基本輸入、輸出數(shù)據(jù)操作,結(jié)合實(shí)例形式分析了Java輸入、輸出數(shù)據(jù)相關(guān)函數(shù)使用技巧與操作注意事項(xiàng),需要的朋友可以參考下
    2020-04-04
  • 基于swing開發(fā)彈幕播放器

    基于swing開發(fā)彈幕播放器

    這篇文章主要為大家詳細(xì)介紹了基于swing實(shí)現(xiàn)彈幕播放器的開發(fā)過程,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2017-06-06
  • java獲取百度網(wǎng)盤真實(shí)下載鏈接的方法

    java獲取百度網(wǎng)盤真實(shí)下載鏈接的方法

    這篇文章主要介紹了java獲取百度網(wǎng)盤真實(shí)下載鏈接的方法,涉及java針對URL操作及頁面分析的相關(guān)技巧,具有一定參考借鑒價(jià)值,需要的朋友可以參考下
    2015-07-07
  • java遞歸算法實(shí)例分析

    java遞歸算法實(shí)例分析

    這篇文章主要介紹了java遞歸算法實(shí)例分析,具有一定借鑒價(jià)值,需要的朋友可以參考下。
    2017-12-12
  • Java字符串編碼知識(shí)點(diǎn)詳解介紹

    Java字符串編碼知識(shí)點(diǎn)詳解介紹

    在本篇內(nèi)容了小編給大家詳細(xì)分析了關(guān)于Java字符串編碼的知識(shí)點(diǎn)并對實(shí)例做了分析,有興趣的朋友們跟著學(xué)習(xí)下。
    2022-11-11
  • Java函數(shù)接口和Lambda表達(dá)式深入分析

    Java函數(shù)接口和Lambda表達(dá)式深入分析

    這篇文章主要介紹了Java函數(shù)接口和Lambda表達(dá)式,函數(shù)接口是一個(gè)具有單個(gè)抽象方法的接口,接口設(shè)計(jì)主要是為了支持Lambda表達(dá)式和方法引用,使得Java能更方便地實(shí)現(xiàn)函數(shù)式編程風(fēng)格,需要的朋友可以參考下
    2025-04-04
  • Java中的StringBuilder性能測試

    Java中的StringBuilder性能測試

    這篇文章主要介紹了Java中的StringBuilder性能測試,本文包含測試代碼和測試結(jié)果,最后得出結(jié)論,需要的朋友可以參考下
    2014-09-09
  • 不知道面試會(huì)不會(huì)問Lambda怎么用(推薦)

    不知道面試會(huì)不會(huì)問Lambda怎么用(推薦)

    這篇文章主要介紹了Lambda表達(dá)式用法,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-04-04
  • Java中json格式化BigDecimal保留2位小數(shù)

    Java中json格式化BigDecimal保留2位小數(shù)

    這篇文章主要給大家介紹了關(guān)于Java中json格式化BigDecimal保留2位小數(shù)的相關(guān)資料,BigDecimal是Java中的一個(gè)數(shù)學(xué)庫,可以實(shí)現(xiàn)高精度計(jì)算,文中給出了詳細(xì)的代碼實(shí)例,需要的朋友可以參考下
    2023-09-09
  • 深入理解Java設(shè)計(jì)模式之組合模式

    深入理解Java設(shè)計(jì)模式之組合模式

    這篇文章主要介紹了JAVA設(shè)計(jì)模式之組合模式的的相關(guān)資料,文中示例代碼非常詳細(xì),供大家參考和學(xué)習(xí),感興趣的朋友可以了解下
    2021-11-11

最新評論

安仁县| 东城区| 库车县| 靖西县| 都兰县| 阳泉市| 三门县| 平舆县| 民权县| 常德市| 桐柏县| 芜湖县| 晋江市| 江山市| 衡山县| 奉化市| 泗洪县| 北辰区| 甘德县| 乡城县| 砚山县| 卓尼县| 固镇县| 朝阳县| 怀集县| 永昌县| 孝昌县| 东辽县| 称多县| 信宜市| 上高县| 大城县| 探索| 双柏县| 太保市| 上林县| 通许县| 丹凤县| 丹棱县| 抚顺县| 永济市|