用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方法下載接口文件中文亂碼問題解決,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧2024-05-05
springboot controller 增加指定前綴的兩種實現方法
這篇文章主要介紹了springboot controller 增加指定前綴的兩種實現方法,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2022-02-02
Spring Boot LocalDateTime格式化處理的示例詳解
這篇文章主要介紹了Spring Boot LocalDateTime格式化處理的示例詳解,小編覺得挺不錯的,現在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧2018-10-10
Spring Cache自定義緩存key和過期時間的實現代碼
使用 Redis的客戶端 Spring Cache時,會發(fā)現生成 key中會多出一個冒號,而且有一個空節(jié)點的存在,查看源碼可知,這是因為 Spring Cache默認生成key的策略就是通過兩個冒號來拼接,本文給大家介紹了Spring Cache自定義緩存key和過期時間的實現,需要的朋友可以參考下2024-05-05
spring應用中多次讀取http post方法中的流遇到的問題
這篇文章主要介紹了spring應用中多次讀取http post方法中的流,文中給大家列舉處理問題描述及解決方法,需要的朋友可以參考下2018-11-11

