• 瞅一瞅JUC提供的限流工具Semaphore


            微服务架构发展到今天,各个服务之前的调用、接口请求越来越频繁,服务器承受的压力自然也越来越大。如果放任所有请求请求到服务器,不管是服务器也好还是数据库也好,都可能被因为无法承受大批量的请求而阻塞、宕机甚至是GG,尤其是一些查询接口请求特别频繁的情况下,是需要对请求进行一定限制的。

            譬如:钉钉打卡数据,动辄几十上百万数据,这种数据一秒钟来个几百次请求,钉钉估计也够呛,调用频率肯定有限制。一些支付流水,对账单也是这个原理,在请求方面都有一些限制。

            还有,服务接口限制,频繁点击提示亦或是Tomcat对超出限制的请求进行丢弃等,都是为了保护服务而实行的一些保护措施。

            之前已经说了漏桶算法令牌桶算法滑动时间窗格,这些都是限流实现的一种方式,今天,我们来看一看JUC给我们提供的一个限流工具类——Semaphore!

            

            先写一个简单的Semaphore使用,大家先对Semaphore有个印象,当然比较熟悉的直接略过哈~

    1. public static void main(String[] args) {
    2. //创建Semaphore,需要指定permits,也就是许可证,或者说信号量
    3. Semaphore semaphore = new Semaphore(3);
    4. // 线程一获取两个许可证(信号量)
    5. new Thread(()->{
    6. System.out.println(Thread.currentThread().getName()+"开始执行了");
    7. boolean flag = semaphore.tryAcquire(2);
    8. if (flag){
    9. System.out.println(Thread.currentThread().getName() + "成功获取到两个信号量,休息10s");
    10. try {
    11. TimeUnit.SECONDS.sleep(20);
    12. semaphore.release(2);
    13. } catch (InterruptedException e) {
    14. e.printStackTrace();
    15. }
    16. }
    17. }).start();
    18. // 线程二获取两个许可证(信号量)
    19. new Thread(()->{
    20. try {
    21. System.out.println(Thread.currentThread().getName() + "开始执行了");
    22. int i = semaphore.availablePermits();
    23. System.out.println("剩余信号量 = " + i);
    24. // 许可证不够,会进行等待尝试,直到25S之后。刚开始有3个许可证,上面线程获取到了2个,除非
    25. // sleep之后释放,否则这里无法获取到足够的许可证。有些类似令牌桶
    26. boolean flag = semaphore.tryAcquire(2,25, TimeUnit.SECONDS);
    27. if (flag) {
    28. System.out.println(Thread.currentThread().getName() + "成功获取到两个信号量,休息10s");
    29. TimeUnit.SECONDS.sleep(10);
    30. }
    31. } catch (Exception e){
    32. e.printStackTrace();
    33. }
    34. }).start();
    35. }

            Semaphore的核心就是permits!也就是许可证,或者说信号量!创建Semaphore的时候,必须指定许可证的数量,有点类似令牌桶,就是你先放好令牌,任务来了,先去取令牌,拿到令牌才能执行任务,没令牌就等着或者直接尥蹶子,返回失败~

               

    下面就来看一看Semaphore的源码,看一看大佬Doug Lea是怎么实现限流滴。

    1. //构造方法必须传许可证数量,默认非公平
    2. public Semaphore(int permits) {
    3. sync = new NonfairSync(permits);
    4. }
    5. static final class NonfairSync extends Sync {
    6. // 默认非公平实现
    7. NonfairSync(int permits) {
    8. super(permits);
    9. }
    10. protected int tryAcquireShared(int acquires) {
    11. return nonfairTryAcquireShared(acquires);
    12. }
    13. }
    14. Sync(int permits) {
    15. setState(permits);
    16. }
    17. protected final void setState(int newState) {
    18. state = newState;
    19. }

    我们可以看到,新建Semaphore中传递的permits最终也就是设置state的数值!也就是说这个permits就是AbstractQueuedSynchronizer(AQS)中共享锁的个数!

    接着就是尝试获取信号量了。

    1. // 获取信号量
    2. public boolean tryAcquire(int permits) {
    3. // 获取小于0的信号量,完全没意义,异常
    4. if (permits < 0) throw new IllegalArgumentException();
    5. // 是否获取成功
    6. return sync.nonfairTryAcquireShared(permits) >= 0;
    7. }
    8. // 非公平尝试获取共享锁
    9. final int nonfairTryAcquireShared(int acquires) {
    10. // 无限循环
    11. for (;;) {
    12. // 获取当前剩余的信号量
    13. int available = getState();
    14. int remaining = available - acquires;
    15. // 如果剩余信号量小于零,说明不够,直接返回,||后面部分不执行,上文会
    16. // 根据>=0为成功,小于零自然是失败
    17. if (remaining < 0 ||
    18. // CAS 设置剩余信号量,设置成功返回否则重新循环
    19. compareAndSetState(available, remaining))
    20. return remaining;
    21. }
    22. }

    信号量其实也就是共享锁个数或者说state的数值,尝试去获取就先看下state是否大于需要的permits,小于的话自然是失败。如果大于或者等于,需要使用CAS进行替换,因为这里的remaining缓存到本线程了,会有并发问题,接着,CAS成功就返回,失败就重新执行这个过程。

    当然,这种只是最简单的获取方式,一般是会在某个时间段内获取,超时之后才算失败。也就是semaphore.tryAcquire(2,25, TimeUnit.SECONDS),一起来瞅瞅~

    1. // 尝试获取信号量,过期时间以及单位
    2. public boolean tryAcquire(int permits, long timeout, TimeUnit unit) throws InterruptedException {
    3. //信号量小于0,无意义,异常
    4. if (permits < 0) throw new IllegalArgumentException();
    5. return sync.tryAcquireSharedNanos(permits, unit.toNanos(timeout));
    6. }
    7. // 固定时间内尝试获取信号量
    8. public final boolean tryAcquireSharedNanos(int arg, long nanosTimeout) throws InterruptedException {
    9. if (Thread.interrupted())
    10. throw new InterruptedException();
    11. return tryAcquireShared(arg) >= 0 ||
    12. doAcquireSharedNanos(arg, nanosTimeout);
    13. }
    14. //上文已讲过,不多说
    15. protected int tryAcquireShared(int acquires) {
    16. return nonfairTryAcquireShared(acquires);
    17. }
    18. // 固定时间内获取arg个信号量
    19. private boolean doAcquireSharedNanos(int arg, long nanosTimeout) throws InterruptedException {
    20. // 时间过期,结束返回
    21. if (nanosTimeout <= 0L)
    22. return false;
    23. //当前系统时间加过期时间,获取结束时间点
    24. final long deadline = System.nanoTime() + nanosTimeout;
    25. //新建共享节点,关联当前线程,并添加到同步队列尾部,
    26. //如果队列为空,执行(enq),初始化同步队列,并将节点添加到尾部
    27. final Node node = addWaiter(Node.SHARED);
    28. boolean failed = true;
    29. try {
    30. // 循环
    31. for (;;) {
    32. // 当前节点的前一个节点
    33. final Node p = node.predecessor();
    34. //如果节点是头结点,获取信号量并返回
    35. if (p == head) {
    36. int r = tryAcquireShared(arg);
    37. // 获取成功
    38. if (r >= 0) {
    39. // 设置头结点
    40. setHeadAndPropagate(node, r);
    41. p.next = null; // help GC
    42. failed = false;
    43. return true;
    44. }
    45. }
    46. // 剩余时间
    47. nanosTimeout = deadline - System.nanoTime();
    48. // 过期,直接返回失败
    49. if (nanosTimeout <= 0L)
    50. return false;
    51. // 尝试获取信号量失败,清理同步队列中无用或者已经线程结束的节点
    52. // 并且剩余时间要大于1000ns,否则不用中断,继续循环即可,1000ns太短
    53. // 程序执行需要时间,可以直接进行循环处理
    54. if (shouldParkAfterFailedAcquire(p, node) &&
    55. nanosTimeout > spinForTimeoutThreshold)
    56. LockSupport.parkNanos(this, nanosTimeout);
    57. //中断异常
    58. if (Thread.interrupted())
    59. throw new InterruptedException();
    60. }
    61. } finally {
    62. if (failed)
    63. cancelAcquire(node);
    64. }
    65. }

    这个获取只要分两步,如果前置节点就是头节点,也就是说就只有当前线程获取信号量,直接去获取即可。第二步,如果不是头结点,自然要排队等待,节点状态设置为Node.SIGNAL,同时清理同步队列中已经状态waitStatus > 0的节点,然后线程中断,等待唤醒。

    其中还有一个setHeadAndPropagate方法,主要的意思就是,上一个线程释放的信号量,可能当前线程获取之后还有盈余,于是,考虑唤醒下一个等待的线程继续获取信号量。这个方法个人感觉理解也是有点绕,推荐一个博客,想深究的可以去瞅下

    1. //设置队列头,并检查后续队列是否正在等待
    2. private void setHeadAndPropagate(Node node, int propagate) {
    3. // 当前节点获取到共享锁,成为头节点
    4. Node h = head;
    5. setHead(node);
    6. //propagate > 0 说明有剩余的信号量,后续节点可以尝试获取信号量,故要doReleaseShared
    7. //propagate == 0 且h.waitStatus < 0当时tryAcquireShared后没有共享锁剩余,但之后的时刻很可能又有共享锁释放出来了。
    8. if (propagate > 0 || h == null || h.waitStatus < 0 ||
    9. (h = head) == null || h.waitStatus < 0) {
    10. Node s = node.next;
    11. if (s == null || s.isShared())
    12. doReleaseShared();
    13. }
    14. }

    接着就是信号量释放了,也就是相当于令牌桶,走完通道,令牌肯定要放回桶里呗

    1. //信号量释放
    2. public void release(int permits) {
    3. // 释放小于0,完全没意义,异常
    4. if (permits < 0) throw new IllegalArgumentException();
    5. // 释放信号量
    6. sync.releaseShared(permits);
    7. }
    8. // 释放信号量(共享锁)
    9. public final boolean releaseShared(int arg) {
    10. if (tryReleaseShared(arg)) {
    11. doReleaseShared();
    12. return true;
    13. }
    14. return false;
    15. }
    16. // Semaphore实现
    17. protected final boolean tryReleaseShared(int releases) {
    18. for (;;) {
    19. int current = getState();
    20. int next = current + releases;
    21. // 要么,release是负数,要么就是超出位数限制,导致比当前还小,
    22. // 异常:超过最大允许计数
    23. if (next < current) // overflow
    24. throw new Error("Maximum permit count exceeded");
    25. //CAS替换
    26. if (compareAndSetState(current, next))
    27. return true;
    28. }
    29. }

    释放信号量的过程也比较简单,就是state重新加上要释放的信号量,然后进行CAS替换即可,当然,还要走一个必须的过程!doReleaseShared!

    1. //释放共享锁
    2. private void doReleaseShared() {
    3. //循环
    4. for (;;) {
    5. Node h = head;
    6. //至少有两个node,就一个节点,也不用释放了
    7. if (h != null && h != tail) {
    8. //获取头结点状态
    9. int ws = h.waitStatus;
    10. // 头结点是SIGNAL,CAS唤醒
    11. if (ws == Node.SIGNAL) {
    12. if (!compareAndSetWaitStatus(h, Node.SIGNAL, 0))
    13. continue; // loop to recheck cases
    14. unparkSuccessor(h);
    15. }
    16. //如果状态为0,说明h的后继所代表的线程已经被唤醒或即将被唤醒,并且这个中间状态即将消失,要么由于acquire thread获取锁失败再次设置head为SIGNAL并再次阻塞,要么由于acquire thread获取锁成功而将自己(head后继)设置为新head并且只要head后继不是队尾,那么新head肯定为SIGNAL。所以设置这种中间状态的head的status为PROPAGATE,让其status又变成负数,这样可能被被唤醒线程检测到
    17. else if (ws == 0 &&
    18. !compareAndSetWaitStatus(h, 0, Node.PROPAGATE))
    19. continue; // loop on failed CAS
    20. }
    21. // head不变退出循环,不然会执行多次循环。
    22. if (h == head) // loop if head changed
    23. break;
    24. }
    25. }

    Semaphore的主要逻辑就是这样,如果只是简单的尝试获取信号量的话,就是直接修改state!不过一般这种肯定是会有一个尝试的过期时间(timeout)的,也就是semaphore.tryAcquire(2,25, TimeUnit.SECONDS);

            这个时候,就需要用到CLH队列了,也就是阻塞等待,由于每次尝试获取的信号量不一致,可能一个线程释放6个信号量,然后后面三个线程,每个都获取两个,线程就会在获取的时候一次被唤醒,这里的逻辑看起来就比较绕,最好是本地造一些数据,debug看下,这样应该能更好的 理解Semaphore中信号量的获取与释放。

    好了,到这里就game over了,感觉还行的话,点个赞呗~

            

    no sacrifice,no victory~

  • 相关阅读:
    Java从零到就业一站通关,解决你的担忧
    第六章:Java内存模型之JMM
    Docker基本命令
    Vue全局事件总线实现任意组件间通信
    nginx -s reload, 提示 [emerg] duplicate location “/“
    测试老鸟整理,从手工测试到自动化测试的进阶全程...
    与客户沟通需要注意什么?
    PyQt5中为窗口添加菜单工具栏状态栏
    52_Pandas处理日期和时间列(字符串转换、日期提取等)
    CAS 学习笔记
  • 原文地址:https://blog.csdn.net/zsah2011/article/details/125560023