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

用JAVA實(shí)現(xiàn)一套背壓機(jī)制

 更新時(shí)間:2023年06月30日 08:47:27   作者:hwp0710  
背壓依我的理解來(lái)說(shuō),是指訂閱者能和發(fā)布者交互,可以調(diào)節(jié)發(fā)布者發(fā)布數(shù)據(jù)的速率,解決把訂閱者壓垮的問(wèn)題,這篇文章主要介紹了用JAVA自己實(shí)現(xiàn)一套背壓機(jī)制,需要的朋友可以參考下

Reactive Streams:一種支持背壓的異步數(shù)據(jù)流處理標(biāo)準(zhǔn),主流實(shí)現(xiàn)有RxJava和Reactor,Spring WebFlux默認(rèn)集成的是Reactor。

Reactive Streams主要解決背壓(back-pressure)問(wèn)題。當(dāng)傳入的任務(wù)速率大于系統(tǒng)處理能力時(shí),數(shù)據(jù)處理將會(huì)對(duì)未處理數(shù)據(jù)產(chǎn)生一個(gè)緩沖區(qū)。

背壓依我的理解來(lái)說(shuō),是指訂閱者能和發(fā)布者交互(通過(guò)代碼里面的調(diào)用request和cancel方法交互),可以調(diào)節(jié)發(fā)布者發(fā)布數(shù)據(jù)的速率,解決把訂閱者壓垮的問(wèn)題。關(guān)鍵在于上面例子里面的訂閱關(guān)系Subscription這個(gè)接口,他有request和cancel 2個(gè)方法,用于通知發(fā)布者需要數(shù)據(jù)和通知發(fā)布者不再接受數(shù)據(jù)。

我們重點(diǎn)理解背壓在jdk9里面是如何實(shí)現(xiàn)的。關(guān)鍵在于發(fā)布者Publisher的實(shí)現(xiàn)類(lèi)SubmissionPublisher的submit方法是block方法。訂閱者會(huì)有一個(gè)緩沖池,默認(rèn)為Flow.defaultBufferSize() = 256。當(dāng)訂閱者的緩沖池滿(mǎn)了之后,發(fā)布者調(diào)用submit方法發(fā)布數(shù)據(jù)就會(huì)被阻塞,發(fā)布者就會(huì)停(慢)下來(lái);訂閱者消費(fèi)了數(shù)據(jù)之后(調(diào)用Subscription.request方法),緩沖池有位置了,submit方法就會(huì)繼續(xù)執(zhí)行下去,就是通過(guò)這樣的機(jī)制,實(shí)現(xiàn)了調(diào)節(jié)發(fā)布者發(fā)布數(shù)據(jù)的速率,消費(fèi)得快,生成就快,消費(fèi)得慢,發(fā)布者就會(huì)被阻塞,當(dāng)然就會(huì)慢下來(lái)了。

單線(xiàn)程版本:

一個(gè)生產(chǎn)者,一個(gè)消費(fèi)者

import lombok.SneakyThrows;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;
public class BackpressureExample {
    public static void main(String[] args) throws InterruptedException {
        BackpressureSubscriber subscriber = new BackpressureSubscriber();
        BackpressurePublisher publisher = new BackpressurePublisher(subscriber);
        publisher.start();
        subscriber.start();
        // 為了演示效果,這里讓主線(xiàn)程休眠一段時(shí)間
        Thread.sleep(50000);
        publisher.stop();
        subscriber.stop();
    }
    @SneakyThrows
    public static void processDataLogic(List<Integer> batch) {
        //模擬任務(wù)執(zhí)行
        int r = new Random().nextInt(3000);
        Thread.sleep(r);
        System.out.println(Thread.currentThread().getName() + ",Received batch: " + batch + ",sleep ms = " + r);
    }
    static class BackpressurePublisher {
        private final BackpressureSubscriber subscriber;
        private volatile boolean running;
        public BackpressurePublisher(BackpressureSubscriber subscriber) {
            this.subscriber = subscriber;
            this.running = true;
        }
        public void start() {
            Thread thread = new Thread(() -> {
                int item = 1;
                while (running) {
                    List<Integer> batch = new ArrayList<>();
                    for (int i = 0; i < 5; i++) {
                        System.out.println(Thread.currentThread().getName() + "-----produce data = " + item);
                        batch.add(item++);
                    }
                    while (!subscriber.accept(batch)) {
                        if (!running) {
                            break;
                        }
                    }
                }
            });
            thread.start();
        }
        public void stop() {
            running = false;
        }
    }
    static class BackpressureSubscriber {
        private volatile boolean running;
        public BackpressureSubscriber() {
            this.running = true;
        }
        public boolean accept(List<Integer> batch) {
            if (running) {
                processDataLogic(batch);
                return true;
            } else {
                return false;
            }
        }
        public void start() {
            // Subscriber 在 JDK 8 中沒(méi)有異步處理的能力,因此不需要單獨(dú)開(kāi)啟線(xiàn)程
        }
        public void stop() {
            running = false;
        }
    }
}

多線(xiàn)程版本

一個(gè)生產(chǎn)者,多個(gè)消費(fèi)者

import lombok.SneakyThrows;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
public class BackpressureExample {
    public static void main(String[] args) throws InterruptedException {
        BackpressureSubscriber subscriber = new BackpressureSubscriber();
        BackpressurePublisher publisher = new BackpressurePublisher(subscriber);
        publisher.start();
        subscriber.start();
        // 為了演示效果,這里讓主線(xiàn)程休眠一段時(shí)間
        Thread.sleep(50000);
        publisher.stop();
        subscriber.stop();
    }
    @SneakyThrows
    public static void processDataLogic(List<Integer> batch) {
        //模擬任務(wù)執(zhí)行
        int r = new Random().nextInt(3000);
        Thread.sleep(r);
        System.out.println(Thread.currentThread().getName() + ",Received batch: " + batch + ",sleep ms = " + r);
    }
    static class BackpressurePublisher {
        private final BackpressureSubscriber subscriber;
        private volatile boolean running;
        public BackpressurePublisher(BackpressureSubscriber subscriber) {
            this.subscriber = subscriber;
            this.running = true;
        }
        public void start() {
            Thread thread = new Thread(() -> {
                int item = 1;
                while (running) {
                    List<Integer> batch = new ArrayList<>();
                    for (int i = 0; i < 5; i++) {
                        System.out.println(Thread.currentThread().getName() + "-----produce data = " + item);
                        batch.add(item++);
                    }
                    while (!subscriber.accept(batch)) {
                        if (!running) {
                            break;
                        }
                    }
                }
            });
            thread.start();
        }
        public void stop() {
            running = false;
        }
    }
    static class BackpressureSubscriber {
        private volatile boolean running;
        private final ExecutorService executor;
        private final int workerSize = 2;
        private final List<Future> futures;
        public BackpressureSubscriber() {
            this.running = true;
            this.executor = Executors.newFixedThreadPool(workerSize);
            futures = new ArrayList<>(workerSize);
        }
        public boolean accept(List<Integer> batch) {
            if (running) {
                Future f = executor.submit(() -> processDataLogic(batch));
                futures.add(f);
                waitForTaskDone(futures);
                return true;
            } else {
                return false;
            }
        }
        public void waitForTaskDone(List<Future> futures) {
            while (futures.size() >= workerSize) {
                for (Future future : futures) {
                    if (future.isDone()) {
                        // 只要有一個(gè)worker是空閑就重新獲取任務(wù)
                        futures.remove(future);
                        return;
                    }
                }
            }
        }
        public void start() {
            // Subscriber 在 JDK 8 中沒(méi)有異步處理的能力,因此不需要單獨(dú)開(kāi)啟線(xiàn)程
        }
        public void stop() {
            running = false;
            executor.shutdown();
        }
    }
}

到此這篇關(guān)于用JAVA自己實(shí)現(xiàn)一套背壓機(jī)制的文章就介紹到這了,更多相關(guān)java背壓機(jī)制內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 使用ServletUtil.write方法下載接口文件中文亂碼問(wèn)題解決

    使用ServletUtil.write方法下載接口文件中文亂碼問(wèn)題解決

    本文主要介紹了使用ServletUtil.write方法下載接口文件中文亂碼問(wèn)題解決,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2024-05-05
  • Java中使用BigDecimal進(jìn)行精確運(yùn)算

    Java中使用BigDecimal進(jìn)行精確運(yùn)算

    這篇文章主要介紹了Java中使用BigDecimal進(jìn)行精確運(yùn)算的方法,非常不錯(cuò),需要的朋友參考下
    2017-02-02
  • Java獲取用戶(hù)訪(fǎng)問(wèn)IP及地理位置的方法詳解

    Java獲取用戶(hù)訪(fǎng)問(wèn)IP及地理位置的方法詳解

    這篇文章主要介紹了Java獲取用戶(hù)訪(fǎng)問(wèn)IP及地理位置的方法,結(jié)合實(shí)例形式詳細(xì)分析了Java基于百度地圖開(kāi)放平臺(tái)獲取用戶(hù)訪(fǎng)問(wèn)IP及地理位置相關(guān)操作技巧,需要的朋友可以參考下
    2020-04-04
  • Spring中配置數(shù)據(jù)源的幾種方式

    Spring中配置數(shù)據(jù)源的幾種方式

    今天小編就為大家分享一篇關(guān)于Spring中配置數(shù)據(jù)源的幾種方式,小編覺(jué)得內(nèi)容挺不錯(cuò)的,現(xiàn)在分享給大家,具有很好的參考價(jià)值,需要的朋友一起跟隨小編來(lái)看看吧
    2019-01-01
  • springboot controller 增加指定前綴的兩種實(shí)現(xiàn)方法

    springboot controller 增加指定前綴的兩種實(shí)現(xiàn)方法

    這篇文章主要介紹了springboot controller 增加指定前綴的兩種實(shí)現(xiàn)方法,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-02-02
  • Spring Boot LocalDateTime格式化處理的示例詳解

    Spring Boot LocalDateTime格式化處理的示例詳解

    這篇文章主要介紹了Spring Boot LocalDateTime格式化處理的示例詳解,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2018-10-10
  • Spring 處理 HTTP 請(qǐng)求參數(shù)注解的操作方法

    Spring 處理 HTTP 請(qǐng)求參數(shù)注解的操作方法

    這篇文章主要介紹了Spring 處理 HTTP 請(qǐng)求參數(shù)注解的操作方法,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),感興趣的朋友參考下吧
    2024-04-04
  • Spring Cache自定義緩存key和過(guò)期時(shí)間的實(shí)現(xiàn)代碼

    Spring Cache自定義緩存key和過(guò)期時(shí)間的實(shí)現(xiàn)代碼

    使用 Redis的客戶(hù)端 Spring Cache時(shí),會(huì)發(fā)現(xiàn)生成 key中會(huì)多出一個(gè)冒號(hào),而且有一個(gè)空節(jié)點(diǎn)的存在,查看源碼可知,這是因?yàn)?nbsp;Spring Cache默認(rèn)生成key的策略就是通過(guò)兩個(gè)冒號(hào)來(lái)拼接,本文給大家介紹了Spring Cache自定義緩存key和過(guò)期時(shí)間的實(shí)現(xiàn),需要的朋友可以參考下
    2024-05-05
  • Java使用Callable接口實(shí)現(xiàn)多線(xiàn)程的實(shí)例代碼

    Java使用Callable接口實(shí)現(xiàn)多線(xiàn)程的實(shí)例代碼

    這篇文章主要介紹了Java使用Callable接口實(shí)現(xiàn)多線(xiàn)程的實(shí)例代碼,實(shí)現(xiàn)Callable和實(shí)現(xiàn)Runnable類(lèi)似,但是功能更強(qiáng)大,具體表現(xiàn)在可以在任務(wù)結(jié)束后提供一個(gè)返回值,Runnable不行,call方法可以?huà)伋霎?Runnable的run方法不行,需要的朋友可以參考下
    2023-08-08
  • spring應(yīng)用中多次讀取http post方法中的流遇到的問(wèn)題

    spring應(yīng)用中多次讀取http post方法中的流遇到的問(wèn)題

    這篇文章主要介紹了spring應(yīng)用中多次讀取http post方法中的流,文中給大家列舉處理問(wèn)題描述及解決方法,需要的朋友可以參考下
    2018-11-11

最新評(píng)論

来宾市| 涟水县| 棋牌| 宁城县| 通道| 南昌市| 天柱县| 广饶县| 曲麻莱县| 金平| 来安县| 苍山县| 姚安县| 民和| 且末县| 广西| 当雄县| 富宁县| 昌宁县| 巴南区| 安吉县| 调兵山市| 盖州市| 曲松县| 通许县| 波密县| 金华市| 黔西县| 乌兰浩特市| 广饶县| 宁夏| 婺源县| 曲松县| 阳原县| 南陵县| 天峨县| 股票| 嘉祥县| 阿克| 西青区| 尼勒克县|