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

java通過信號量實現(xiàn)限流的示例

 更新時間:2023年06月29日 09:50:52   作者:Shawn_Shawn  
本文主要介紹了java通過信號量實現(xiàn)限流的示例,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

信號量(Semaphore)是 Java 多線程并發(fā)中的一種 JDK 內(nèi)置同步器,通過它可以實現(xiàn)多線程對公共資源的并發(fā)訪問控制。

信號量由來

限流器

信號量的主要應(yīng)用場景是控制最多 N 個線程同時地訪問資源,其中計數(shù)器的最大值即是許可的最大值 N。

現(xiàn)在我們需要開發(fā)一個限流器,同一時刻最多有10個請求可以執(zhí)行。對于這樣的需求,我們實現(xiàn)的方案有:

  • 使用Atomic類
  • 使用Lock
  • 使用條件變量
  • 使用信號量

使用Atomic類實現(xiàn)

public class LimitByAtomic {
?
 ?private static final AtomicInteger COUNTER = new AtomicInteger(10);
?
 ?public void f() {
 ? ?int count = COUNTER.decrementAndGet();
 ? ?if (count < 0) {
 ? ? ?COUNTER.incrementAndGet();
 ? ? ?System.out.println("拒絕執(zhí)行業(yè)務(wù)邏輯");
 ? ? ?return; // 拒絕執(zhí)行業(yè)務(wù)邏輯
 ?  }
?
 ? ?try {
 ? ? ?// 執(zhí)行業(yè)務(wù)邏輯
 ? ? ?System.out.println("執(zhí)行業(yè)務(wù)邏輯");
 ?  } finally {
 ? ? ?COUNTER.incrementAndGet();
 ?  }
  }
}

使用Lock實現(xiàn)

public class LimitByLock {
?
 ?private int count = 10;
?
 ?public void f() {
 ? ?if (count <= 0) {
 ? ? ?System.out.println("拒絕執(zhí)行業(yè)務(wù)邏輯");
 ? ? ?return;
 ?  }
?
 ? ?synchronized (this) {
 ? ? ?if (count <= 0) {
 ? ? ? ?System.out.println("拒絕執(zhí)行業(yè)務(wù)邏輯");
 ? ? ? ?return;
 ? ?  }
 ? ? ?count--;
 ?  }
?
 ? ?try {
 ? ? ?// 執(zhí)行業(yè)務(wù)邏輯
 ? ? ?System.out.println("執(zhí)行業(yè)務(wù)邏輯");
 ?  } finally {
 ? ? ?synchronized (this) {
 ? ? ? ?count++;
 ? ?  }
 ?  }
  }
}

使用條件變量實現(xiàn)

對于使用Atomic類還是Lock這兩種實現(xiàn)方式,都有一個缺點,如果10個線程同時執(zhí)行,當?shù)?1個線程來執(zhí)行的時候,會被拒絕掉,這樣就沒有執(zhí)行業(yè)務(wù)邏輯的機會,造成請求丟失。

所以我們可以通過線程等待-通知機制來解決上面的問題。如果10個線程同時執(zhí)行,當?shù)?1個線程來執(zhí)行的時候,先阻塞這第11個線程,等待前面的10個線程只要執(zhí)行完一個,就通知第11個線程來執(zhí)行。

public class LimitByCondition {
?
 ?private int count = 10;
?
 ?public void f() throws Exception {
 ? ?synchronized (this) {
 ? ? ?while (count <= 0) {
 ? ? ? ?System.out.println("等待執(zhí)行業(yè)務(wù)邏輯");
 ? ? ? ?this.wait();
 ? ?  }
 ? ? ?count--;
 ?  }
?
 ? ?try {
 ? ? ?System.out.println("執(zhí)行業(yè)務(wù)邏輯");
 ?  } finally {
 ? ? ?synchronized (this) {
 ? ? ? ?count++;
 ? ? ? ?this.notifyAll();
 ? ?  }
 ?  }
  }
}

使用Semaphore實現(xiàn)

除了使用條件變量,java sdk中還可以使用Semaphore來實現(xiàn)。

public class LimitBySemaphore {
?
 ?private final Semaphore semaphore = new Semaphore(10);
?
 ?public void f() throws Exception {
 ? ?semaphore.acquire();
 ? ?try {
 ? ? ?System.out.println("執(zhí)行業(yè)務(wù)邏輯");
 ?  } finally {
 ? ? ?semaphore.release();
 ?  }
  }
}

接下來我們就來探討一下Semaphore的實現(xiàn)原理。

Semaphore實現(xiàn)原理

信號量模型

實際上Semaphore的實現(xiàn)原理非常簡單,總結(jié)下來就是:一個計數(shù)器,一個等待隊列,三個方法。

在信號量模型里,計數(shù)器和等待隊列對外是透明,所以只能通過信號量模型提供的三個方法訪問,init(),down(),up()--這些方法都是原子性的。

init():設(shè)置計數(shù)器的初始值。

down():計數(shù)器的值減1;如果此時計數(shù)器的值小于0,則當前線程將被阻塞,否則當前線程可以繼續(xù)執(zhí)行。

up():計數(shù)器的值加1;如果此時計數(shù)器的值小于等于0,則喚醒等待隊列中的一個線程,并將其從等待隊列中移除。

class MySemaphore{
 ?// 計數(shù)器
 ?int count;
 ?// 等待隊列
 ?Queue queue;
 ?// 初始化操作
 ?MySemaphore(int c){
 ? ?this.count=c;
  }
 ?// 
 ?void down(){
 ? ?this.count--;
 ? ?if(this.count<0){
 ? ? ?// 將當前線程插入等待隊列
 ? ? ?// 阻塞當前線程
 ?  }
  }
 ?void up(){
 ? ?this.count++;
 ? ?if(this.count<=0) {
 ? ? ?// 移除等待隊列中的某個線程 T
 ? ? ?// 喚醒線程 T
 ?  }
  }
}
?

使用方法如下:

static int count;
// 初始化信號量
static final MySemaphore s 
 ? ?= new MySemaphore(1);
// 用信號量保證互斥 ? ?
static void addOne() {
 ?s.down();
 ?try {
 ? ?count+=1;
  } finally {
 ? ?s.up();
  }
}
?

實際上信號量模型,down()、up() 這兩個操作歷史上最早稱為 P 操作和 V 操作,所以信號量模型也被稱為 PV 原語。

Java Semaphore的實現(xiàn)

public class Semaphore implements java.io.Serializable {
  public void acquire() throws InterruptedException;
 ?public void acquireUninterruptibly();
 ?public boolean tryAcquire();
 ?public boolean tryAcquire(long timeout, TimeUnit unit);
 ?public void release();
 ?public void acquire(int permits) throws InterruptedException;
 ?public void acquireUninterruptibly(int permits) ;
 ?public boolean tryAcquire(int permits);
 ?public boolean tryAcquire(int permits, long timeout, TimeUnit unit)
 ? ?throws InterruptedException;
 ?public void release(int permits);
}

Java Semaphore的實現(xiàn),acquire()對應(yīng)信號量模型里的down()方法,release()對應(yīng)信號量模型里的up()方法。

Semaphore類提供的常用方法有以下幾個。我們可以粗略地將以下方法分為兩組。前五個為一組,默認一次獲取或釋放的許可(permit)個數(shù)為1。后五個為一組,可以指定一次獲取或釋放的許可個數(shù)。對于每組方法來說,都有4個不同的獲取許可的方法:可中斷獲取、不可中斷獲取、非阻塞獲取、可超時獲取,這跟Lock提供的各種加鎖方法非常相似。

Java Semaphore的實現(xiàn)也是基于AQS來實現(xiàn)的,跟ReentrantLock一樣,Semaphore中的AQS也有公平鎖與非公平鎖這兩種實現(xiàn)。

public class Semaphore implements java.io.Serializable {
    abstract static class Sync extends AbstractQueuedSynchronizer {
 ?  }
?
 ? ?// 非公平鎖
 ? ?static final class NonfairSync extends Sync {
 ?  }
 ? ?// 公平鎖
 ? ?static final class FairSync extends Sync {
 ?  }
 ? ?// 默認使用非公平鎖
 ? ?public Semaphore(int permits) {
 ? ? ? ?sync = new NonfairSync(permits);
 ?  }
?
 ? ?public Semaphore(int permits, boolean fair) {
 ? ? ? ?sync = fair ? new FairSync(permits) : new NonfairSync(permits);
 ?  }
}

Semaphore可以看做是一種共享鎖,因此,F(xiàn)airSync類和NofairSync類實現(xiàn)了AQS的tryAcquireShared()抽象方法,不過,實現(xiàn)邏輯并不相同。對于tryReleaseShared()抽象方法,因為在FairSync和NofairSync中的實現(xiàn)邏輯相同,因此,它被放置于FairSync和NofairSync的公共父類Sync中。

acquire()實現(xiàn)如下:

// java.util.concurrent.Semaphore#acquire()
public void acquire() throws InterruptedException {
 ?sync.acquireSharedInterruptibly(1);
}
?
// java.util.concurrent.locks.AbstractQueuedSynchronizer#acquireSharedInterruptibly
public final void acquireSharedInterruptibly(int arg)
 ? ? ? ? ? ?throws InterruptedException {
 ?// 先判斷線程有沒有被中斷
 ?if (Thread.interrupted())
 ? ?throw new InterruptedException();
 ?// 嘗試獲取共享鎖,如果獲取許可失敗,返回值<0, 需要進入等待隊列
 ?if (tryAcquireShared(arg) < 0)
 ? ?doAcquireSharedInterruptibly(arg); // 排隊等待隊列
}
?

tryAcquireShared()實現(xiàn)

// java.util.concurrent.Semaphore.FairSync#tryAcquireShared
protected int tryAcquireShared(int acquires) {
 ?for (;;) {
 ? ?if (hasQueuedPredecessors()) // 比非公平鎖多了這一行
 ? ? ?return -1;
 ? ?int available = getState();
 ? ?int remaining = available - acquires;
 ? ?if (remaining < 0 ||
 ? ? ? ?compareAndSetState(available, remaining))
 ? ? ?return remaining;
  }
}
?
// java.util.concurrent.Semaphore.Sync#nonfairTryAcquireShared
final int nonfairTryAcquireShared(int acquires) {
 ?for (;;) {
 ? ?int available = getState();
 ? ?int remaining = available - acquires;
 ? ?if (remaining < 0 ||
 ? ? ? ?compareAndSetState(available, remaining))
 ? ? ?return remaining;
  }
}

以上兩個tryAcquireShared()函數(shù)的代碼實現(xiàn)基本相同。許可個數(shù)存放在AQS的state變量中,兩個函數(shù)都是通過自旋+CAS的方式來獲取許可。兩個函數(shù)唯一的區(qū)別在于,對于公平模式下的Semaphore,當線程調(diào)用tryAcquireShared()函數(shù)時,如果等待隊列中有等待許可的線程,那么,線程將直接去排隊等待許可,而不是像非公平模式下的Semaphore那樣,線程可以插隊直接競爭許可。

release()實現(xiàn)

// java.util.concurrent.Semaphore#release()
public void release() {
 ?sync.releaseShared(1);
}
?
// java.util.concurrent.locks.AbstractQueuedSynchronizer#releaseShared
public final boolean releaseShared(int arg) {
 ?// 嘗試釋放許可
 ?if (tryReleaseShared(arg)) {
 ? ?// 喚醒等待隊列其中一個線程
 ? ?doReleaseShared();
 ? ?return true;
  }
 ?return false;
}
?
// java.util.concurrent.Semaphore.Sync#tryReleaseShared
protected final boolean tryReleaseShared(int releases) {
 ?// 采用自旋 + CAS來更新state
 ?for (;;) {
 ? ?int current = getState();
 ? ?int next = current + releases;
 ? ?if (next < current) // overflow
 ? ? ?throw new Error("Maximum permit count exceeded");
 ? ?if (compareAndSetState(current, next))
 ? ? ?return true;
  }
}

總結(jié)

semaphore其中一個功能是lock不容易實現(xiàn)的,那就是:semaphore可以允許多個線程訪問同一個臨界區(qū)。

比較常見的需求就是我們工作中遇到各種池化資源,例如連接池,對象池,線程池等等。其中,最熟悉的可能是數(shù)據(jù)庫連接池,在同一時刻,一定是允許多個線程同時使用連接池的,當然,每個鏈接在被釋放前,是不允許其他線程使用的。

對象池:一次性創(chuàng)建出N個對象,之后所有的線程重復利用這N個對象,對象在被釋放前,也是不允許其他線程使用的。對象池,可以用List保存實例對象。

class ObjPool<T, R> {
 ?final List<T> pool;
 ?// 用信號量實現(xiàn)限流器
 ?final Semaphore sem;
 ?// 構(gòu)造函數(shù)
 ?ObjPool(int size, T t){
 ? ?pool = new Vector<T>(){};
 ? ?for(int i=0; i<size; i++){
 ? ? ?pool.add(t);
 ?  }
 ? ?sem = new Semaphore(size);
  }
 ?// 利用對象池的對象,調(diào)用 func 限流
 ?R exec(Function<T,R> func) {
 ? ?T t = null;
 ? ?sem.acquire();
 ? ?try {
 ? ? ?t = pool.remove(0);
 ? ? ?return func.apply(t);
 ?  } finally {
 ? ? ?pool.add(t);
 ? ? ?sem.release();
 ?  }
  }
}
// 創(chuàng)建對象池
ObjPool<Long, String> pool = 
 ?new ObjPool<Long, String>(10, 2);
// 通過對象池獲取 t,之后執(zhí)行 ?
pool.exec(t -> {
 ? ?System.out.println(t);
 ? ?return t.toString();
});
?

到此這篇關(guān)于java通過信號量實現(xiàn)限流的示例的文章就介紹到這了,更多相關(guān)java 限流內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 基于mybatis-plus-generator實現(xiàn)代碼自動生成器

    基于mybatis-plus-generator實現(xiàn)代碼自動生成器

    這篇文章專門為小白準備了入門級mybatis-plus-generator代碼自動生成器,可以提高開發(fā)效率。文中的示例代碼講解詳細,感興趣的可以了解一下
    2022-05-05
  • Java中的異步操作CompletableFuture示例詳解

    Java中的異步操作CompletableFuture示例詳解

    CompletableFuture是Java8引入的異步編程工具,支持非阻塞的鏈式調(diào)用和組合操作,提供了強大的異步任務(wù)編排能力,本文通過實例代碼講解Java中的異步操作CompletableFuture,感興趣的朋友跟隨小編一起看看吧
    2026-02-02
  • Java使用鎖解決銀行取錢問題實例分析

    Java使用鎖解決銀行取錢問題實例分析

    這篇文章主要介紹了Java使用鎖解決銀行取錢問題,結(jié)合實例形式分析了java線程同步與鎖機制相關(guān)原理及操作注意事項,需要的朋友可以參考下
    2019-08-08
  • SpringBoot使用自定義注解+AOP+Redis實現(xiàn)接口限流的實例代碼

    SpringBoot使用自定義注解+AOP+Redis實現(xiàn)接口限流的實例代碼

    這篇文章主要介紹了SpringBoot使用自定義注解+AOP+Redis實現(xiàn)接口限流,本文通過實例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-09-09
  • Spring事件監(jiān)聽機制使用和原理解析

    Spring事件監(jiān)聽機制使用和原理解析

    Spring的監(jiān)聽機制基于觀察者模式,就是就是我們所說的發(fā)布訂閱模式,這種模式可以在一定程度上實現(xiàn)代碼的解耦,本文將從原理上解析Spring事件監(jiān)聽機制,需要的朋友可以參考下
    2023-06-06
  • SpringCloud使用Nacos 配置中心實現(xiàn)配置自動刷新功能使用

    SpringCloud使用Nacos 配置中心實現(xiàn)配置自動刷新功能使用

    SpringCloud項目中使用Nacos作為配置中心可以方便開發(fā)及運維人員隨時查看配置信息,及配置共享,并且Nacos支持配置信息實時刷新,非常方便,下面給大家介紹SpringCloud使用Nacos配置中心實現(xiàn)配置自動刷新功能使用,感興趣的朋友一起看看吧
    2025-05-05
  • java實現(xiàn)對excel文件的處理合并單元格的操作

    java實現(xiàn)對excel文件的處理合并單元格的操作

    這篇文章主要介紹了java實現(xiàn)對excel文件的處理合并單元格的操作,開頭給大家介紹了依賴引入代碼,表格操作的核心代碼,代碼超級簡單,需要的朋友可以參考下
    2021-07-07
  • Java并發(fā)編程之線程池實現(xiàn)原理詳解

    Java并發(fā)編程之線程池實現(xiàn)原理詳解

    池化思想是一種空間換時間的思想,期望使用預先創(chuàng)建好的對象來減少頻繁創(chuàng)建對象的性能開銷,java中有多種池化思想的應(yīng)用,例如:數(shù)據(jù)庫連接池、線程池等,下面就來具體講講
    2023-05-05
  • mybatis防止SQL注入的方法實例詳解

    mybatis防止SQL注入的方法實例詳解

    SQL注入是一種很簡單的攻擊手段,但直到今天仍然十分常見。那么mybatis是如何防止SQL注入的呢?下面腳本之家小編給大家?guī)砹藢嵗a,需要的朋友參考下吧
    2018-04-04
  • mybatis-plus在pom文件中的出錯問題及解決

    mybatis-plus在pom文件中的出錯問題及解決

    文章內(nèi)容為提示更新MyBatis-Plus框架至最新版本,并需同步更新Spring?Boot依賴,建議補全Maven倉庫中MyBatis-Plus的artifact?ID(如com.baomidou:mybatis-plus-boot-starter),并參考官方倉庫獲取準確版本信息
    2025-08-08

最新評論

菏泽市| 自贡市| 汕尾市| 道孚县| 荣昌县| 南漳县| 土默特左旗| 伊宁县| 兴宁市| 英吉沙县| 阿拉善右旗| 仁布县| 西峡县| 五家渠市| 荆州市| 苗栗市| 凭祥市| 白银市| 西平县| 尉氏县| 湟源县| 阿鲁科尔沁旗| 驻马店市| 临邑县| 通州市| 临汾市| 沿河| 通许县| 太和县| 原阳县| 金乡县| 全州县| 金溪县| 泰安市| 五常市| 宁晋县| 开远市| 诏安县| 潜江市| 根河市| 甘洛县|