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

Project?Reactor源碼解析publishOn使用示例

 更新時間:2022年08月15日 17:05:15   作者:夜盡天明_  
這篇文章主要為大家介紹了Project?Reactor源碼解析publishOn使用示例,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪

功能分析

相關(guān)示例源碼:github.com/chentianmin…

public final Flux<T> publishOn(Scheduler scheduler, boolean delayError, int prefetch)

onNext()onComplete()、onError()方法進行線程切換,publishOn()使得它下游的消費階段異步執(zhí)行。

  • scheduler:線程切換的調(diào)度器,Scheduler用來生成實際執(zhí)行異步任務(wù)的Worker
  • delayError:是否延時轉(zhuǎn)發(fā)Error。如果為true,當(dāng)收到上游的Error時,會等隊列中的元素消費完畢后再向下游轉(zhuǎn)發(fā)Error。否則會立即轉(zhuǎn)發(fā)Error,可能導(dǎo)致隊列中的元素丟失。默認(rèn)為true。
  • prefetch:預(yù)取元素的數(shù)量,同時也是隊列的容量。默認(rèn)值為Queues.SMALL_BUFFER_SIZE,該值通過配置進行修改。

代碼示例

prefetch

/**
 * 每隔delayMillis生產(chǎn)一個元素
 */
protected Flux<Integer> delayPublishFlux(int delayMillis, int startInclusive, int endExclusive) {
    return Flux.create(fluxSink -> {
        IntStream.range(startInclusive, endExclusive)
                .forEach(i -> {
                    // 同步next
                    sleep(delayMillis);
                    logInt(i, "生產(chǎn)");
                    fluxSink.next(i);
                });
        fluxSink.complete();
    });
}
@Test
public void testPreFetch() {
    delayPublishFlux(1000, 1, 5)
            .doOnRequest(i -> logLong(i, "request"))
            .publishOn(Schedulers.boundedElastic(), 2)
            .subscribe(i -> logInt(i, "消費"));
    sleep(10000);
}

每次會都向上游請求2個元素。另外還能發(fā)現(xiàn),從第二個request開始,線程發(fā)生了切換。

delayError

/**
 * 每隔delayMillis生產(chǎn)一個元素,最后發(fā)送Error
 */
protected Flux<Integer> delayPublishFluxError(int delayMillis, int startInclusive, int endExclusive) {
    return Flux.create(fluxSink -> {
        IntStream.range(startInclusive, endExclusive)
                .forEach(i -> {
                    // 同步next
                    sleep(delayMillis);
                    logInt(i, "生產(chǎn)");
                    fluxSink.next(i);
                });
        fluxSink.error(new RuntimeException("發(fā)布錯誤!"));
    });
}
@Test
public void testDelayError() {
    delayPublishFluxError(500, 1, 5)
            .publishOn(Schedulers.boundedElastic())
            // 只是為了消費慢一點
            .doOnNext(i -> sleep(1000))
            .subscribe(i -> logInt(i, "消費"));
    sleep(10000);
}

元素消費完才觸發(fā)Error!

@Test
public void testNotDelayError() {
    delayPublishFluxError(500, 1, 5)
            .publishOn(Schedulers.boundedElastic(), false, 256)
            // 只是為了消費慢一點
            .doOnNext(i -> sleep(1000))
            .subscribe(i -> logInt(i, "消費"));
    sleep(10000);
}

元素還沒消費完就觸發(fā)Error!

源碼分析

首先看一下publishOn()操作符在裝配階段做了什么,直接查看Flux#publishOn()源碼。

Flux#publishOn()

publishOn()裝配階段重點是創(chuàng)建了FluxPublishOn對象。

接下來,我們分析訂閱階段發(fā)生了什么。一個Publisher在訂閱的時候調(diào)用的是其subscribe()方法,因此我們繼續(xù)看Flux#subscribe()源碼。

Flux#subscribe()

Flux#subscribe()方法的實現(xiàn)中,如果上游PublisherOptimizableOperator類型,實際的Subscriber是通過調(diào)用該InternalFluxOperator#subscribeOrReturn()方法返回的。如果返回值為null,直接return。

對于publishOn()操作符來說,裝配階段創(chuàng)建的FluxPublishOn就是OptimizableOperator類型。所以繼續(xù)查看FluxPublishOn#subscribeOrReturn()源碼。

FluxPublishOn#subscribeOrReturn()

可以看到,方法返回的是PublishOnSubscriber,它包裝了原始的Subscriber

在后續(xù)的訂閱階段一定會調(diào)用其onSubscribe()方法,在運行階段一定會調(diào)用其onNext()方法。我們先看FluxPublishOn#onSubscribe()源碼。

FluxPublishOn#onSubscribe()

onSubscribe()實現(xiàn)中,分為同步隊列融合、異步隊列融合以及非融合方式處理。

如果上游的SubscriptionQueueSubscription類型,則會進行隊列融合。具體采用同步還是異步,取決于該QueueSubscription#requestFusion()實現(xiàn)。

  • 同步隊列融合:復(fù)用當(dāng)前隊列,繼續(xù)調(diào)用下游onSubscribe()方法,但不會繼續(xù)調(diào)用上游request()方法。
  • 異步隊列融合:復(fù)用當(dāng)前隊列,然后繼續(xù)調(diào)用下游onSubscribe()以及上游request()方法,請求數(shù)量是prefetch。
  • 非融合:創(chuàng)建一個新的隊列,然后繼續(xù)調(diào)用下游onSubscribe()以及上游request()方法,請求數(shù)量是prefetch。

接下來,我們從源碼角度分別介紹上述三種方式的處理邏輯,首先介紹非融合方式。

非融合

先看如下代碼示例,該代碼會以非融合方式執(zhí)行。

@Test
public void testNoFuse() {
    delayPublishFlux(1000, 1, 5)
            .publishOn(Schedulers.boundedElastic())
            .subscribe(i -> logInt(i, "消費"));
    sleep(10000);
}

間隔1s生產(chǎn)消費元素!

在消費階段,一定會調(diào)用FluxPublishOn#onNext()方法。

FluxPublishOn#onNext()

我們重點關(guān)注非融合方式執(zhí)行邏輯,其實只做了2件事:

  • 將下發(fā)的元素添加到隊列中,該隊列就是onSubscribe()階段創(chuàng)建的新隊列。
  • 調(diào)用trySchedule()方法進行調(diào)度。

繼續(xù)看FluxPublishOn#trySchedule()源碼。

FluxPublishOn#trySchedule()

這里其實就是交由woker異步執(zhí)行,后續(xù)會執(zhí)行FluxPublishOn.run()方法。

FluxPublishOn#run()

在run()方法執(zhí)行的時候,分為3段邏輯:

  • 如果是輸出融合,執(zhí)行runBackfused()方法。
  • 如果是同步隊列融合,執(zhí)行runSync()方法。
  • 否則,執(zhí)行runAsync()方法。

對于當(dāng)前例子,實際執(zhí)行的是runAsync()方法,繼續(xù)查看其源碼。

FluxPublishOn#runAsync()

runAsync()做的事情比較簡單,就是排空隊列中的元素下發(fā)給下游。同時在這里會繼續(xù)調(diào)用request()向上游請求數(shù)據(jù),這也是前面說的從第二個request()開始會進行線程切換的原因。

另外這里還會調(diào)用checkTerminated(),檢查終止情況。

FluxPublishOn#checkTerminated()

如果delayError=true,必須當(dāng)前隊列為空是才會轉(zhuǎn)發(fā)Error。如果delayError=false,則直接轉(zhuǎn)發(fā)Error。繼續(xù)查看onComplete()方法。

FluxPublishOn#onComplete()

如果未結(jié)束,將done標(biāo)記設(shè)置為true,然后再次調(diào)用trySchedule()進行調(diào)度。后續(xù)再被調(diào)度到的時候,如果隊列已經(jīng)排空,才會調(diào)用下游onComplete(),觸發(fā)完成。

小結(jié)

簡單總結(jié)一下非融合執(zhí)行過程:

onSubscribe()時創(chuàng)建一個隊列,在onNext()時將上游下發(fā)的元素添加到隊列中,然后異步排空隊列中的元素,繼續(xù)下發(fā)給下游。

同步隊列融合

以下代碼會以同步隊列融合方式執(zhí)行。

@Test
public void testSyncFuse() {
    Flux.just(1, 2 ,3, 4, 5)
            .publishOn(Schedulers.boundedElastic())
            .subscribe(this::logInt);
    sleep(10000);
}

因為Flux.just()對應(yīng)的SubscriptionSynchronousSubscription,其requestFusion()方法實現(xiàn)如下:

SynchronousSubscription#requestFusion()

此時返回的是SYNC,執(zhí)行同步隊列融合。

前面提到過,同步隊列融合會復(fù)用當(dāng)前隊列,繼續(xù)調(diào)用下游onSubscribe()方法,但不會繼續(xù)調(diào)用上游request()方法。

這意味著,此時FluxPublishOn#onNext()FluxPublishOn#onComplete()方法并不會調(diào)用。但是FluxPublishOn#request()依然會被下游調(diào)用到。

FluxPublishOn#request()

request()方法中還是會調(diào)用trySchedule(),后續(xù)會異步調(diào)用runSync()方法(前面已經(jīng)分析了)。

對于非融合方式,trySchedule()也會執(zhí)行,只是這次調(diào)度的時候,隊列中還沒有數(shù)據(jù)被添加進去。

FluxPublishOn#runSync()

runSync()實現(xiàn)上runAsync()差不多,也是排空隊列的元素,繼續(xù)下發(fā)給下游。不同的點是少了request()調(diào)用,以及取消完成控制有差異。

小結(jié)

簡單總結(jié)一下同步隊列融合執(zhí)行過程:

onSubsrribe()時直接復(fù)用上游QueueSubscription作為隊列,不會調(diào)用上游request()請求數(shù)據(jù),在自身request()時異步排空隊列中的元素,繼續(xù)下發(fā)給下游。

異步隊列融合

以下代碼會以異步隊列融合方式執(zhí)行。

@Test
public void testAsyncFuse() {
    Flux.just(1, 2, 3, 4, 5)
            .windowUntil(i -&gt; i % 3 == 0)
            .publishOn(Schedulers.boundedElastic())
            .flatMap(Function.identity())
            .subscribe(this::logInt);
    sleep(10000);
}

因為windowUntil()對應(yīng)的SubscriptionWindowPredicateMain,其requestFusion()方法實現(xiàn)如下:

WindowPredicateMain#requestFusion()

此時返回ASYNC,執(zhí)行異步隊列融合。接下來再看一下FluxPublishOn#onNext()源碼。

FluxPublishOn#onNext()

注意,此時onNext()方法參數(shù)是null,表明上游并沒有真正下發(fā)元素,可以將其看做是一個觸發(fā)Worker調(diào)度的信號。后續(xù)還是會異步執(zhí)行runAsync()方法,這里就不再分析了。

這其實也很容易理解:異步隊列融合直接復(fù)用了上游的QueueSubscription作為隊列,真正的數(shù)據(jù)應(yīng)該由這個隊列下發(fā)。

總結(jié)

簡單總結(jié)一下同步隊列融合執(zhí)行過程:

onSubsrribe()時直接復(fù)用上游QueueSubscription作為隊列,在onNext()時接收上游信號,異步排空隊列中的元素,繼續(xù)下發(fā)給下游。

非融合、同步隊列融合、異步隊列融合比較如下:

以上就是Project Reactor源碼解析publishOn使用示例的詳細內(nèi)容,更多關(guān)于Project Reactor publishOn的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • java圖片壓縮工具類

    java圖片壓縮工具類

    這篇文章主要為大家詳細介紹了java圖片壓縮工具類,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-02-02
  • JPA like 模糊查詢 語法格式解析

    JPA like 模糊查詢 語法格式解析

    這篇文章主要介紹了JPA like 模糊查詢 語法格式解析,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-12-12
  • JavaWeb實現(xiàn)用戶登錄注冊功能實例代碼(基于Servlet+JSP+JavaBean模式)

    JavaWeb實現(xiàn)用戶登錄注冊功能實例代碼(基于Servlet+JSP+JavaBean模式)

    這篇文章主要基于Servlet+JSP+JavaBean開發(fā)模式實現(xiàn)JavaWeb用戶登錄注冊功能實例代碼,非常實用,本文介紹的非常詳細,具有參考借鑒價值,感興趣的朋友一起看看吧
    2016-05-05
  • Netty粘包問題的常見解決方案

    Netty粘包問題的常見解決方案

    粘包和拆包問題也叫做粘包和半包問題,它是指在數(shù)據(jù)傳輸時,接收方未能正常讀取到一條完整數(shù)據(jù)的情況(只讀取了部分?jǐn)?shù)據(jù),或多讀取到了另一條數(shù)據(jù)的情況)就叫做粘包或拆包問題,本文介紹了Netty如何解決粘包問題,需要的朋友可以參考下
    2024-06-06
  • JetBrains IntelliJ IDEA 優(yōu)化教超詳細程

    JetBrains IntelliJ IDEA 優(yōu)化教超詳細程

    這篇文章主要介紹了JetBrains IntelliJ IDEA 優(yōu)化教超詳細程,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-03-03
  • SpringBoot實現(xiàn)列表數(shù)據(jù)導(dǎo)出為Excel文件

    SpringBoot實現(xiàn)列表數(shù)據(jù)導(dǎo)出為Excel文件

    這篇文章主要為大家詳細介紹了在Spring?Boot框架中如何將列表數(shù)據(jù)導(dǎo)出為Excel文件,文中的示例代碼講解詳細,感興趣的小伙伴可以了解下
    2024-02-02
  • mybatis-plus 實現(xiàn)分頁查詢的示例代碼

    mybatis-plus 實現(xiàn)分頁查詢的示例代碼

    本文介紹了在MyBatis-Plus中實現(xiàn)分頁查詢,包括引入依賴、配置分頁插件、使用分頁查詢以及在控制器中調(diào)用分頁查詢的方法,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2024-11-11
  • Eclipse智能提示及快捷鍵

    Eclipse智能提示及快捷鍵

    本文主要介紹了Eclipse智能提示及快捷鍵的相關(guān)知識,具有很好的參考價值。下面跟著小編一起來看下吧
    2017-03-03
  • 關(guān)于jpa?querydsl嵌套查詢demo

    關(guān)于jpa?querydsl嵌套查詢demo

    這篇文章主要介紹了關(guān)于jpa?querydsl?嵌套查詢demo,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-05-05
  • iOS多線程介紹

    iOS多線程介紹

    這篇文章主要介紹了iOS多線程的相關(guān)知識,涉及到對進程,線程等方面的知識講解,本文非常具有參考價值,感興趣的朋友一起學(xué)習(xí)吧
    2016-05-05

最新評論

清原| 砀山县| 于都县| 福贡县| 崇礼县| 南汇区| 阜南县| 泰和县| 云霄县| 新郑市| 樟树市| 青浦区| 沙田区| 芷江| 屯昌县| 崇仁县| 丰都县| 大荔县| 千阳县| 城步| 江达县| 彭山县| 栖霞市| 迁安市| 安多县| 金山区| 措美县| 页游| 广汉市| 宣武区| 磐安县| 东乌| 安远县| 山东| 大洼县| 铜鼓县| 乐山市| 平罗县| 鱼台县| 保靖县| 丹巴县|