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

JAVA?Reactor中Sinks.Many類三種常見(jiàn)的創(chuàng)建方式及使用

 更新時(shí)間:2025年07月29日 10:32:41   作者:?jiǎn)鑶琛? 
Sinks.Many用于創(chuàng)建多值(Multi-Value)的發(fā)布者(Publisher)的一種機(jī)制,它允許用戶將數(shù)據(jù)從一個(gè)地方發(fā)送到多個(gè)訂閱者,這篇文章主要介紹了JAVA?Reactor中Sinks.Many類三種常見(jiàn)的創(chuàng)建方式及使用,需要的朋友可以參考下

一、Java中的Sinks.Many

在Java編程中,我們經(jīng)常需要處理數(shù)據(jù)流。為了更有效地處理數(shù)據(jù)流,我們可以使用Reactor庫(kù)中的Sinks.Many類。這個(gè)類提供了一種簡(jiǎn)單而強(qiáng)大的方式來(lái)處理多個(gè)事件流,并且可以通過(guò)異步或者同步的方式處理這些事件。

Sinks.Many是Reactor庫(kù)中的一個(gè)類,它可以用來(lái)處理多個(gè)事件流。它提供了一種簡(jiǎn)單而強(qiáng)大的方式來(lái)處理多個(gè)事件,可以在任何時(shí)間添加、獲取以及完成事件。

Sinks.Many為什么被設(shè)計(jì)為 Sinks 類的內(nèi)部接口?

實(shí)現(xiàn)細(xì)節(jié)隱藏??: Sinks.Many 的實(shí)際實(shí)現(xiàn)(如 MulticastReplayProcessor 或其他內(nèi)部類)被封裝在 Sinks 類的內(nèi)部。用戶只需通過(guò)工廠方法(如 Sinks.many().multicast())獲取接口,無(wú)需關(guān)心底層實(shí)現(xiàn)。 ??

接口職責(zé)分離??: Sinks 類作為工廠,統(tǒng)一管理所有類型的 Sink(如 Sinks.One、Sinks.Many),而 Sinks.Many 作為內(nèi)部接口,明確表示它是一個(gè)??多播數(shù)據(jù)源??,與單播(Sinks.One)等場(chǎng)景隔離。

二、Sinks.Many的創(chuàng)建

在源碼中可以看待,Sinks.Many提供了三中常見(jiàn)的創(chuàng)建方式。

(1)unicast()

Sinks.Many<String> unicastSink = Sinks.many().unicast().onBackpressureBuffer();

這種創(chuàng)建方式中提供設(shè)置背壓緩沖區(qū)的方法

  • 用途

  • 創(chuàng)建一個(gè) 單播(Unicast) 的 Sinks.Many,僅允許 一個(gè)訂閱者 訂閱。

  • 核心特性

    • 單訂閱者限制:第二個(gè)訂閱者嘗試訂閱時(shí)會(huì)觸發(fā) IllegalStateException。

    • 背壓支持:通過(guò)緩沖區(qū)處理生產(chǎn)者和消費(fèi)者的速率不匹配,默認(rèn)使用無(wú)界緩沖區(qū)(需手動(dòng)配置限制)。

    • 無(wú)歷史數(shù)據(jù):訂閱者只能收到訂閱后產(chǎn)生的數(shù)據(jù)。

  • 適用場(chǎng)景

    • 點(diǎn)對(duì)點(diǎn)通信(如任務(wù)隊(duì)列)。

    • 需要嚴(yán)格保證單訂閱者的場(chǎng)景(如資源獨(dú)占型操作)。

(2)multicast()

Sinks.Many<String> multicastSink = Sinks.many().multicast().onBackpressureBuffer(100);

這個(gè)創(chuàng)建方式,相比較Unicast,就是可以允許多個(gè)訂閱者訂閱,同樣提供了多種設(shè)置緩存區(qū)的方式

  • 用途
    創(chuàng)建一個(gè) 多播(Multicast) 的 Sinks.Many,允許多個(gè)訂閱者訂閱。

  • 核心特性

    • 多訂閱者支持:所有訂閱者共享同一數(shù)據(jù)流。

    • 無(wú)歷史數(shù)據(jù):新訂閱者只能收到訂閱后產(chǎn)生的數(shù)據(jù),無(wú)法獲取之前的數(shù)據(jù)。

    • 背壓策略:默認(rèn)使用無(wú)界緩沖區(qū),但可以通過(guò)配置限制(如 onBackpressureBuffer(int capacity))。

  • 適用場(chǎng)景

    • 實(shí)時(shí)廣播(如股票行情推送)。

    • 需要多個(gè)消費(fèi)者并行處理相同數(shù)據(jù)的場(chǎng)景。

(3)replay()

Sinks.Many<String> replaySink = Sinks.many().replay().all();
這個(gè)方法中提供了限制存放歷史數(shù)據(jù)的方法

  • 用途
    創(chuàng)建一個(gè) 支持?jǐn)?shù)據(jù)重放(Replay) 的 Sinks.Many,允許多個(gè)訂閱者訂閱,并回放歷史數(shù)據(jù)。

  • 核心特性

    • 歷史數(shù)據(jù)緩存:新訂閱者可以收到訂閱前一定數(shù)量的數(shù)據(jù)(通過(guò) limit(int) 或 time(Duration) 配置)。

    • 多訂閱者支持:與 multicast() 類似,但增加了數(shù)據(jù)重放能力。

    • 內(nèi)存管理:緩存數(shù)據(jù)量或時(shí)間窗口可配置,避免內(nèi)存無(wú)限增長(zhǎng)。

  • 適用場(chǎng)景

    • 需要新訂閱者獲取歷史數(shù)據(jù)的場(chǎng)景(如聊天記錄回放)。

    • 實(shí)時(shí)監(jiān)控面板(多個(gè)訂閱者需要看到完整上下文)。

三、Sinks.Many的使用

(1)添加事件

一旦我們創(chuàng)建了一個(gè)Sinks.Many對(duì)象,我們可以使用emitNext()方法來(lái)添加一個(gè)事件。這個(gè)方法接受一個(gè)參數(shù),表示要添加的事件。

Sinks.Many<String> sink = Sinks.many().unicast().onBackpressureBuffer();
sink.emitNext("Event 1");
sink.emitNext("Event 2");
sink.emitNext("Event 3");

我們首先創(chuàng)建了一個(gè)Sinks.Many對(duì)象,然后使用emitNext()方法添加了三個(gè)事件。

當(dāng)然,如果我們選擇接受Flux流中的數(shù)據(jù)的時(shí)候,可以這樣添加數(shù)據(jù)

 //創(chuàng)建Flux流
        Flux<String> flux = Flux.<String>create(sink->{
            sink.next("你好");
            sink.complete();
        });

        //創(chuàng)建Sinks.Many處理流數(shù)據(jù)
        Sinks.Many<String> sink = Sinks.many().unicast().onBackpressureBuffer();

        //訂閱Flux流,并將數(shù)據(jù)交給sink處理(像做流數(shù)據(jù)的緩存,篩選流數(shù)據(jù))
        flux.subscribe(
                sink::tryEmitNext,
                sink::tryEmitError,
                sink::tryEmitComplete
        );

(2)獲取流數(shù)據(jù)

我們可以使用Sinks.Many對(duì)象來(lái)獲取已添加的事件??梢允褂胊sFlux()方法將Sinks.Many對(duì)象轉(zhuǎn)換為一個(gè)Flux對(duì)象,然后使用Flux對(duì)象的方法來(lái)訂閱和處理事件。

 // 1. .unicast()使用創(chuàng)建單播 Sink
        Sinks.Many<String> unicastSink = Sinks.many().unicast().onBackpressureBuffer();

        // 2. 將 Sink 轉(zhuǎn)換為 Flux(供訂閱者訂閱)
        Flux<String> flux = unicastSink.asFlux();

        // 3. 第一個(gè)訂閱者(合法)
        flux.subscribe(
                data -> System.out.println("訂閱者1收到數(shù)據(jù): " + data),
                error -> System.err.println("訂閱者1發(fā)生錯(cuò)誤: " + error),
                () -> System.out.println("訂閱者1完成")
        );

        // 4. 推送數(shù)據(jù)到 Sink
        unicastSink.tryEmitNext("數(shù)據(jù)1");
        unicastSink.tryEmitNext("數(shù)據(jù)2");

        // 5. 嘗試第二個(gè)訂閱者(會(huì)拋出 IllegalStateException)
        try {
            flux.subscribe(
                    data -> System.out.println("訂閱者2收到數(shù)據(jù): " + data),
                    error -> System.err.println("訂閱者2發(fā)生錯(cuò)誤: " + error),
                    () -> System.out.println("訂閱者2完成")
            );
        } catch (IllegalStateException e) {
            System.err.println("訂閱者2訂閱失敗: " + e.getMessage());
        }

        // 6. 關(guān)閉 Sink(發(fā)送完成信號(hào))
        unicastSink.tryEmitComplete();

這里我使用了unicast()方法創(chuàng)建的Sinks.Many,這個(gè)時(shí)候我通過(guò)asFlux()方法轉(zhuǎn)換的flux流,只能被一個(gè)訂閱者訂閱到,第二個(gè)訂閱者,訂閱的時(shí)候就出報(bào)錯(cuò),當(dāng)然,如果你想要多個(gè)訂閱者訂閱,可以使用multicast()或者replay()方式創(chuàng)建,

四、熱冷流

上面三種方式創(chuàng)建的是熱流。

熱流:數(shù)據(jù)獨(dú)立于訂閱者持續(xù)生成,多訂閱者共享實(shí)時(shí)數(shù)據(jù),適用于實(shí)時(shí)事件推送。

在添加事件代碼中我也有使用到Flux.create()方法創(chuàng)建流,要注意這里創(chuàng)建的是冷流。

冷流:每次調(diào)用 subscribe() 時(shí),會(huì)觸發(fā) Flux.create() 的回調(diào)函數(shù)(即 Consumer<SynchronousSink<T>>),??重新生成數(shù)據(jù)流??。多個(gè)訂閱者之間數(shù)據(jù)獨(dú)立。這個(gè)跟ThreadLocal有點(diǎn)像,F(xiàn)lux.create()會(huì)為每一個(gè)訂閱者創(chuàng)建單獨(dú)隔離的數(shù)據(jù)流,保證每一條流中數(shù)據(jù)互不影響。

冷熱流的選擇;

  • 若需要多個(gè)訂閱者共享實(shí)時(shí)數(shù)據(jù) → 熱流。
  • 若需要每個(gè)訂閱者獨(dú)立消費(fèi)完整數(shù)據(jù) → 冷流。
  • 若需要?dú)v史數(shù)據(jù) → 使用 replay() 緩存。

總結(jié) 

到此這篇關(guān)于JAVA Reactor中Sinks.Many類三種常見(jiàn)的創(chuàng)建方式及使用的文章就介紹到這了,更多相關(guān)JAVA Reactor中Sinks.Many類內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • SpringBoot實(shí)現(xiàn)系統(tǒng)監(jiān)控的完整示例

    SpringBoot實(shí)現(xiàn)系統(tǒng)監(jiān)控的完整示例

    本文介紹了如何對(duì)SpringBoot應(yīng)用進(jìn)行系統(tǒng)監(jiān)控,包括監(jiān)控的重要性、詳細(xì)步驟以及最佳實(shí)踐,監(jiān)控可以幫助我們預(yù)防問(wèn)題、快速定位和解決,提升系統(tǒng)性能和開(kāi)發(fā)體驗(yàn),需要的朋友可以參考下
    2025-12-12
  • Netty分布式高性能工具類FastThreadLocal和Recycler分析

    Netty分布式高性能工具類FastThreadLocal和Recycler分析

    這篇文章主要為大家介紹了Netty分布式高性能工具類FastThreadLocal和Recycler分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-03-03
  • Maven安裝、配置、Idea配置maven全過(guò)程

    Maven安裝、配置、Idea配置maven全過(guò)程

    這篇文章主要介紹了Maven安裝、配置、Idea配置maven全過(guò)程,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2026-05-05
  • Java線程讓步y(tǒng)ield用法實(shí)例分析

    Java線程讓步y(tǒng)ield用法實(shí)例分析

    這篇文章主要介紹了Java線程讓步y(tǒng)ield用法,結(jié)合實(shí)例形式分析了java中yield()方法的功能、原理及線程讓步操作的相關(guān)實(shí)現(xiàn)技巧,需要的朋友可以參考下
    2019-09-09
  • Linux服務(wù)器Java進(jìn)程消失問(wèn)題解決

    Linux服務(wù)器Java進(jìn)程消失問(wèn)題解決

    這篇文章主要介紹了Linux服務(wù)器Java進(jìn)程消失問(wèn)題解決,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-11-11
  • MybatisPlus中如何調(diào)用Oracle存儲(chǔ)過(guò)程

    MybatisPlus中如何調(diào)用Oracle存儲(chǔ)過(guò)程

    這篇文章主要介紹了MybatisPlus中如何調(diào)用Oracle存儲(chǔ)過(guò)程的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-05-05
  • Spring MVC使用視圖解析的問(wèn)題解讀

    Spring MVC使用視圖解析的問(wèn)題解讀

    這篇文章主要介紹了Spring MVC使用視圖解析的問(wèn)題解讀,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2025-03-03
  • JavaStream將List轉(zhuǎn)為Map示例

    JavaStream將List轉(zhuǎn)為Map示例

    這篇文章主要為大家介紹了JavaStream將List轉(zhuǎn)為Map示例,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-09-09
  • Java中Jackson的多態(tài)反序列化詳解

    Java中Jackson的多態(tài)反序列化詳解

    這篇文章主要介紹了Java中Jackson的多態(tài)反序列化詳解,多態(tài)序列化與反序列化,主要是借助于Jackson的@JsonTypeInfo與@JsonSubTypes注解實(shí)現(xiàn),下面將通過(guò)幾個(gè)例子來(lái)簡(jiǎn)述其運(yùn)用,需要的朋友可以參考下
    2023-11-11
  • Spring?Boot?實(shí)現(xiàn)Redis分布式鎖原理

    Spring?Boot?實(shí)現(xiàn)Redis分布式鎖原理

    這篇文章主要介紹了Spring?Boot實(shí)現(xiàn)Redis分布式鎖原理,文章圍繞主題展開(kāi)詳細(xì)的內(nèi)容介紹,具有一定的參考價(jià)值,需要的朋友可以參考一下
    2022-08-08

最新評(píng)論

兴义市| 华亭县| 洛隆县| 原阳县| 普宁市| 平乐县| 乐平市| 长葛市| 福清市| 池州市| 南漳县| 景德镇市| 防城港市| 独山县| 大足县| 富顺县| 德州市| 巢湖市| 双辽市| 陇西县| 永靖县| 竹溪县| 阿鲁科尔沁旗| 昂仁县| 彝良县| 韶关市| 湖口县| 呼玛县| 米林县| 石台县| 永寿县| 杭州市| 石泉县| 乌拉特中旗| 屏东县| 射洪县| 根河市| 广东省| 囊谦县| 安丘市| 康马县|