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

用JAVA實現一套背壓機制

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

Reactive Streams:一種支持背壓的異步數據流處理標準,主流實現有RxJava和Reactor,Spring WebFlux默認集成的是Reactor。

Reactive Streams主要解決背壓(back-pressure)問題。當傳入的任務速率大于系統(tǒng)處理能力時,數據處理將會對未處理數據產生一個緩沖區(qū)。

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

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

單線程版本:

一個生產者,一個消費者

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();
        // 為了演示效果,這里讓主線程休眠一段時間
        Thread.sleep(50000);
        publisher.stop();
        subscriber.stop();
    }
    @SneakyThrows
    public static void processDataLogic(List<Integer> batch) {
        //模擬任務執(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 中沒有異步處理的能力,因此不需要單獨開啟線程
        }
        public void stop() {
            running = false;
        }
    }
}

多線程版本

一個生產者,多個消費者

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();
        // 為了演示效果,這里讓主線程休眠一段時間
        Thread.sleep(50000);
        publisher.stop();
        subscriber.stop();
    }
    @SneakyThrows
    public static void processDataLogic(List<Integer> batch) {
        //模擬任務執(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()) {
                        // 只要有一個worker是空閑就重新獲取任務
                        futures.remove(future);
                        return;
                    }
                }
            }
        }
        public void start() {
            // Subscriber 在 JDK 8 中沒有異步處理的能力,因此不需要單獨開啟線程
        }
        public void stop() {
            running = false;
            executor.shutdown();
        }
    }
}

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

相關文章

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

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

    本文主要介紹了使用ServletUtil.write方法下載接口文件中文亂碼問題解決,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2024-05-05
  • Java中使用BigDecimal進行精確運算

    Java中使用BigDecimal進行精確運算

    這篇文章主要介紹了Java中使用BigDecimal進行精確運算的方法,非常不錯,需要的朋友參考下
    2017-02-02
  • Java獲取用戶訪問IP及地理位置的方法詳解

    Java獲取用戶訪問IP及地理位置的方法詳解

    這篇文章主要介紹了Java獲取用戶訪問IP及地理位置的方法,結合實例形式詳細分析了Java基于百度地圖開放平臺獲取用戶訪問IP及地理位置相關操作技巧,需要的朋友可以參考下
    2020-04-04
  • Spring中配置數據源的幾種方式

    Spring中配置數據源的幾種方式

    今天小編就為大家分享一篇關于Spring中配置數據源的幾種方式,小編覺得內容挺不錯的,現在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2019-01-01
  • springboot controller 增加指定前綴的兩種實現方法

    springboot controller 增加指定前綴的兩種實現方法

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

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

    這篇文章主要介紹了Spring Boot LocalDateTime格式化處理的示例詳解,小編覺得挺不錯的,現在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-10-10
  • Spring 處理 HTTP 請求參數注解的操作方法

    Spring 處理 HTTP 請求參數注解的操作方法

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

    Spring Cache自定義緩存key和過期時間的實現代碼

    使用 Redis的客戶端 Spring Cache時,會發(fā)現生成 key中會多出一個冒號,而且有一個空節(jié)點的存在,查看源碼可知,這是因為 Spring Cache默認生成key的策略就是通過兩個冒號來拼接,本文給大家介紹了Spring Cache自定義緩存key和過期時間的實現,需要的朋友可以參考下
    2024-05-05
  • Java使用Callable接口實現多線程的實例代碼

    Java使用Callable接口實現多線程的實例代碼

    這篇文章主要介紹了Java使用Callable接口實現多線程的實例代碼,實現Callable和實現Runnable類似,但是功能更強大,具體表現在可以在任務結束后提供一個返回值,Runnable不行,call方法可以拋出異,Runnable的run方法不行,需要的朋友可以參考下
    2023-08-08
  • spring應用中多次讀取http post方法中的流遇到的問題

    spring應用中多次讀取http post方法中的流遇到的問題

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

最新評論

鄂伦春自治旗| 汾西县| 屯昌县| 襄垣县| 久治县| 汾西县| 田东县| 冷水江市| 迁西县| 岢岚县| 马关县| 贞丰县| 秦皇岛市| 清原| 平凉市| 岳阳县| 财经| 渝北区| 杭州市| 江北区| 城市| 盐源县| 昌平区| 谢通门县| 徐汇区| 广德县| 昆山市| 广南县| 汕尾市| 库尔勒市| 澄江县| 桂东县| 高青县| 高雄县| 辽宁省| 高尔夫| 云安县| 旬阳县| 涞源县| 富蕴县| 历史|