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

Netty分布式NioEventLoop任務(wù)隊(duì)列執(zhí)行源碼分析

 更新時(shí)間:2022年03月25日 15:23:02   作者:向南是個(gè)萬人迷  
這篇文章主要為大家介紹了Netty分布式NioEventLoop任務(wù)隊(duì)列執(zhí)行源碼分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

前文傳送門:NioEventLoop處理IO事件

執(zhí)行任務(wù)隊(duì)列

繼續(xù)回到NioEventLoop的run()方法:

protected void run() {
    for (;;) {
        try {
            switch (selectStrategy.calculateStrategy(selectNowSupplier, hasTasks())) {
                case SelectStrategy.CONTINUE:
                    continue;
                case SelectStrategy.SELECT:
                    //輪詢io事件(1)
                    select(wakenUp.getAndSet(false));
                    if (wakenUp.get()) {
                        selector.wakeup();
                    }
                default:
            }
            cancelledKeys = 0;
            needsToSelectAgain = false;
            //默認(rèn)是50
            final int ioRatio = this.ioRatio; 
            if (ioRatio == 100) {
                try {
                    processSelectedKeys();
                } finally {
                    runAllTasks();
                }
            } else {
                //記錄下開始時(shí)間
                final long ioStartTime = System.nanoTime();
                try {
                    //處理輪詢到的key(2)
                    processSelectedKeys();
                } finally {
                    //計(jì)算耗時(shí)
                    final long ioTime = System.nanoTime() - ioStartTime;
                    //執(zhí)行task(3)
                    runAllTasks(ioTime * (100 - ioRatio) / ioRatio);
                }
            }
        } catch (Throwable t) {
            handleLoopException(t);
        }
        //代碼省略
    }
}

我們看到處理完輪詢到的key之后, 首先記錄下耗時(shí), 然后通過runAllTasks(ioTime * (100 - ioRatio) / ioRatio)執(zhí)行taskQueue中的任務(wù)

我們知道ioRatio默認(rèn)是50, 所以執(zhí)行完ioTime * (100 - ioRatio) / ioRatio后, 方法傳入的值為ioTime, 也就是processSelectedKeys()的執(zhí)行時(shí)間:

跟進(jìn)runAllTasks方法:

protected boolean runAllTasks(long timeoutNanos) {
    //定時(shí)任務(wù)隊(duì)列中聚合任務(wù)
    fetchFromScheduledTaskQueue();
    //從普通taskQ里面拿一個(gè)任務(wù)
    Runnable task = pollTask();
    //task為空, 則直接返回
    if (task == null) {
        //跑完所有的任務(wù)執(zhí)行收尾的操作
        afterRunningAllTasks();
        return false;
    }
    //如果隊(duì)列不為空
    //首先算一個(gè)截止時(shí)間(+50毫秒, 因?yàn)閳?zhí)行任務(wù), 不要超過這個(gè)時(shí)間)
    final long deadline = ScheduledFutureTask.nanoTime() + timeoutNanos;
    long runTasks = 0;
    long lastExecutionTime;
    //執(zhí)行每一個(gè)任務(wù)
    for (;;) {
        safeExecute(task);
        //標(biāo)記當(dāng)前跑完的任務(wù)
        runTasks ++;
        //當(dāng)跑完64個(gè)任務(wù)的時(shí)候, 會(huì)計(jì)算一下當(dāng)前時(shí)間
        if ((runTasks & 0x3F) == 0) {
            //定時(shí)任務(wù)初始化到當(dāng)前的時(shí)間
            lastExecutionTime = ScheduledFutureTask.nanoTime();
            //如果超過截止時(shí)間則不執(zhí)行(nanoTime()是耗時(shí)的)
            if (lastExecutionTime >= deadline) {
                break;
            }
        }
        //如果沒有超過這個(gè)時(shí)間, 則繼續(xù)從普通任務(wù)隊(duì)列拿任務(wù)
        task = pollTask();
        //直到?jīng)]有任務(wù)執(zhí)行
        if (task == null) {
            //記錄下最后執(zhí)行時(shí)間
            lastExecutionTime = ScheduledFutureTask.nanoTime();
            break;
        }
    }
    //收尾工作
    afterRunningAllTasks();
    this.lastExecutionTime = lastExecutionTime;
    return true;
}

首先會(huì)執(zhí)行fetchFromScheduledTaskQueue()這個(gè)方法, 這個(gè)方法的意思是從定時(shí)任務(wù)隊(duì)列中聚合任務(wù), 也就是將定時(shí)任務(wù)中找到可以執(zhí)行的任務(wù)添加到taskQueue中

我們跟進(jìn)fetchFromScheduledTaskQueue()方法

private boolean fetchFromScheduledTaskQueue() {
    long nanoTime = AbstractScheduledEventExecutor.nanoTime();
    //從定時(shí)任務(wù)隊(duì)列中抓取第一個(gè)定時(shí)任務(wù)
    //尋找截止時(shí)間為nanoTime的任務(wù)
    Runnable scheduledTask  = pollScheduledTask(nanoTime);
    //如果該定時(shí)任務(wù)隊(duì)列不為空, 則塞到普通任務(wù)隊(duì)列里面
    while (scheduledTask != null) {
        //如果添加到普通任務(wù)隊(duì)列過程中失敗
        if (!taskQueue.offer(scheduledTask)) {
            //則重新添加到定時(shí)任務(wù)隊(duì)列中
            scheduledTaskQueue().add((ScheduledFutureTask<?>) scheduledTask);
            return false;
        }
        //繼續(xù)從定時(shí)任務(wù)隊(duì)列中拉取任務(wù)
        //方法執(zhí)行完成之后, 所有符合運(yùn)行條件的定時(shí)任務(wù)隊(duì)列, 都添加到了普通任務(wù)隊(duì)列中
        scheduledTask = pollScheduledTask(nanoTime);
    }
    return true;
}

 long nanoTime = AbstractScheduledEventExecutor.nanoTime()

 代表從定時(shí)任務(wù)初始化到現(xiàn)在過去了多長(zhǎng)時(shí)間

 Runnable scheduledTask= pollScheduledTask(nanoTime) 

代表從定時(shí)任務(wù)隊(duì)列中拿到小于nanoTime時(shí)間的任務(wù), 因?yàn)樾∮诔跏蓟浆F(xiàn)在的時(shí)間, 說明該任務(wù)需要執(zhí)行了

跟到其父類AbstractScheduledEventExecutor的pollScheduledTask(nanoTime)方法中:

protected final Runnable pollScheduledTask(long nanoTime) {
    assert inEventLoop();
    //拿到定時(shí)任務(wù)隊(duì)列
    Queue<ScheduledFutureTask<?>> scheduledTaskQueue = this.scheduledTaskQueue;
    //peek()方法拿到第一個(gè)任務(wù)
    ScheduledFutureTask<?> scheduledTask = scheduledTaskQueue == null ? null : scheduledTaskQueue.peek();
    if (scheduledTask == null) {
        return null;
    }
    if (scheduledTask.deadlineNanos() <= nanoTime) {
        //從隊(duì)列中刪除
        scheduledTaskQueue.remove();
        //返回該任務(wù)
        return scheduledTask;
    }
    return null;
}

我們看到首先獲得當(dāng)前類綁定的定時(shí)任務(wù)隊(duì)列的成員變量

如果不為空, 則通過scheduledTaskQueue.peek()彈出第一個(gè)任務(wù)

如果當(dāng)前任務(wù)小于傳來的時(shí)間, 說明該任務(wù)需要執(zhí)行, 則從定時(shí)任務(wù)隊(duì)列中刪除

我們繼續(xù)回到fetchFromScheduledTaskQueue()方法中:

private boolean fetchFromScheduledTaskQueue() {
    long nanoTime = AbstractScheduledEventExecutor.nanoTime();
    //從定時(shí)任務(wù)隊(duì)列中抓取第一個(gè)定時(shí)任務(wù)
    //尋找截止時(shí)間為nanoTime的任務(wù)
    Runnable scheduledTask  = pollScheduledTask(nanoTime);
    //如果該定時(shí)任務(wù)隊(duì)列不為空, 則塞到普通任務(wù)隊(duì)列里面
    while (scheduledTask != null) {
        //如果添加到普通任務(wù)隊(duì)列過程中失敗
        if (!taskQueue.offer(scheduledTask)) {
            //則重新添加到定時(shí)任務(wù)隊(duì)列中
            scheduledTaskQueue().add((ScheduledFutureTask<?>) scheduledTask);
            return false;
        }
        //繼續(xù)從定時(shí)任務(wù)隊(duì)列中拉取任務(wù)
        //方法執(zhí)行完成之后, 所有符合運(yùn)行條件的定時(shí)任務(wù)隊(duì)列, 都添加到了普通任務(wù)隊(duì)列中
        scheduledTask = pollScheduledTask(nanoTime);
    }
    return true;
}

彈出需要執(zhí)行的定時(shí)任務(wù)之后, 我們通過taskQueue.offer(scheduledTask)添加到taskQueue中, 如果添加失敗, 則通過

scheduledTaskQueue().add((ScheduledFutureTask<?>) scheduledTask)

重新添加到定時(shí)任務(wù)隊(duì)列中

如果添加成功, 則通過pollScheduledTask(nanoTime)方法繼續(xù)添加, 直到?jīng)]有需要執(zhí)行的任務(wù)

這樣就將定時(shí)任務(wù)隊(duì)列需要執(zhí)行的任務(wù)添加到了taskQueue中

回到runAllTasks(long timeoutNanos)方法中

protected boolean runAllTasks(long timeoutNanos) {
    //定時(shí)任務(wù)隊(duì)列中聚合任務(wù)
    fetchFromScheduledTaskQueue();
    //從普通taskQ里面拿一個(gè)任務(wù)
    Runnable task = pollTask();
    //task為空, 則直接返回
    if (task == null) {
        //跑完所有的任務(wù)執(zhí)行收尾的操作
        afterRunningAllTasks();
        return false;
    }
    //如果隊(duì)列不為空
    //首先算一個(gè)截止時(shí)間(+50毫秒, 因?yàn)閳?zhí)行任務(wù), 不要超過這個(gè)時(shí)間)
    final long deadline = ScheduledFutureTask.nanoTime() + timeoutNanos;
    long runTasks = 0;
    long lastExecutionTime;
    //執(zhí)行每一個(gè)任務(wù)
    for (;;) {
        safeExecute(task);
        //標(biāo)記當(dāng)前跑完的任務(wù)
        runTasks ++;
        //當(dāng)跑完64個(gè)任務(wù)的時(shí)候, 會(huì)計(jì)算一下當(dāng)前時(shí)間
        if ((runTasks & 0x3F) == 0) {
            //定時(shí)任務(wù)初始化到當(dāng)前的時(shí)間
            lastExecutionTime = ScheduledFutureTask.nanoTime();
            //如果超過截止時(shí)間則不執(zhí)行(nanoTime()是耗時(shí)的)
            if (lastExecutionTime >= deadline) {
                break;
            }
        }
        //如果沒有超過這個(gè)時(shí)間, 則繼續(xù)從普通任務(wù)隊(duì)列拿任務(wù)
        task = pollTask();
        //直到?jīng)]有任務(wù)執(zhí)行
        if (task == null) {
            //記錄下最后執(zhí)行時(shí)間
            lastExecutionTime = ScheduledFutureTask.nanoTime();
            break;
        }
    }
    //收尾工作
    afterRunningAllTasks();
    this.lastExecutionTime = lastExecutionTime;
    return true;
}

首先通過 Runnable task = pollTask() 從taskQueue中拿一個(gè)任務(wù)

任務(wù)不為空, 則通過

final long deadline = ScheduledFutureTask.nanoTime() + timeoutNanos 

計(jì)算一個(gè)截止時(shí)間, 任務(wù)的執(zhí)行時(shí)間不能超過這個(gè)時(shí)間

然后在for循環(huán)中通過safeExecute(task)執(zhí)行task

我們跟到safeExecute(task)中:

protected static void safeExecute(Runnable task) {
    try {
        //直接調(diào)用run()方法執(zhí)行
        task.run();
    } catch (Throwable t) {
        //發(fā)生異常不終止
        logger.warn("A task raised an exception. Task: {}", task, t);
    }
}

這里直接調(diào)用task的run()方法進(jìn)行執(zhí)行, 其中發(fā)生異常, 只打印一條日志, 代表發(fā)生異常不終止, 繼續(xù)往下執(zhí)行

回到runAllTasks(long timeoutNanos)方法

protected boolean runAllTasks(long timeoutNanos) {
    //定時(shí)任務(wù)隊(duì)列中聚合任務(wù)
    fetchFromScheduledTaskQueue();
    //從普通taskQ里面拿一個(gè)任務(wù)
    Runnable task = pollTask();
    //task為空, 則直接返回
    if (task == null) {
        //跑完所有的任務(wù)執(zhí)行收尾的操作
        afterRunningAllTasks();
        return false;
    }
    //如果隊(duì)列不為空
    //首先算一個(gè)截止時(shí)間(+50毫秒, 因?yàn)閳?zhí)行任務(wù), 不要超過這個(gè)時(shí)間)
    final long deadline = ScheduledFutureTask.nanoTime() + timeoutNanos;
    long runTasks = 0;
    long lastExecutionTime;
    //執(zhí)行每一個(gè)任務(wù)
    for (;;) {
        safeExecute(task);
        //標(biāo)記當(dāng)前跑完的任務(wù)
        runTasks ++;
        //當(dāng)跑完64個(gè)任務(wù)的時(shí)候, 會(huì)計(jì)算一下當(dāng)前時(shí)間
        if ((runTasks & 0x3F) == 0) {
            //定時(shí)任務(wù)初始化到當(dāng)前的時(shí)間
            lastExecutionTime = ScheduledFutureTask.nanoTime();
            //如果超過截止時(shí)間則不執(zhí)行(nanoTime()是耗時(shí)的)
            if (lastExecutionTime >= deadline) {
                break;
            }
        }
        //如果沒有超過這個(gè)時(shí)間, 則繼續(xù)從普通任務(wù)隊(duì)列拿任務(wù)
        task = pollTask();
        //直到?jīng)]有任務(wù)執(zhí)行
        if (task == null) {
            //記錄下最后執(zhí)行時(shí)間
            lastExecutionTime = ScheduledFutureTask.nanoTime();
            break;
        }
    }
    //收尾工作
    afterRunningAllTasks();
    this.lastExecutionTime = lastExecutionTime;
    return true;
}

每次執(zhí)行完task, runTasks自增

這里 if ((runTasks & 0x3F) == 0) 代表是否執(zhí)行了64個(gè)任務(wù), 如果執(zhí)行了64個(gè)任務(wù), 則會(huì)通過 lastExecutionTime = ScheduledFutureTask.nanoTime() 記錄定時(shí)任務(wù)初始化到現(xiàn)在的時(shí)間, 如果這個(gè)時(shí)間超過了截止時(shí)間, 則退出循環(huán)

如果沒有超過截止時(shí)間, 則通過 task = pollTask() 繼續(xù)彈出任務(wù)執(zhí)行

這里執(zhí)行64個(gè)任務(wù)統(tǒng)計(jì)一次時(shí)間, 而不是每次執(zhí)行任務(wù)都統(tǒng)計(jì), 主要原因是因?yàn)楂@取系統(tǒng)時(shí)間是個(gè)比較耗時(shí)的操作, 這里是netty的一種優(yōu)化方式

如果沒有task需要執(zhí)行, 則通過afterRunningAllTasks()做收尾工作, 最后記錄下最后的執(zhí)行時(shí)間

以上就是有關(guān)執(zhí)行任務(wù)隊(duì)列的相關(guān)邏輯

章節(jié)小結(jié)

本章學(xué)習(xí)了有關(guān)NioEventLoopGroup的創(chuàng)建, NioEventLoop的創(chuàng)建和啟動(dòng), 以及多路復(fù)用器的輪詢處理和task執(zhí)行的相關(guān)邏輯, 通過本章學(xué)習(xí), 我們應(yīng)該掌握如下內(nèi)容:

        1.  NioEventLoopGroup如何選擇分配NioEventLoop

        2.  NioEventLoop如何開啟

        3.  NioEventLoop如何進(jìn)行select操作

        4.  NioEventLoop如何執(zhí)行task

以上就是Netty分布式NioEventLoop任務(wù)隊(duì)列執(zhí)行源碼分析的詳細(xì)內(nèi)容,更多關(guān)于Netty分布式NioEventLoop執(zhí)行任務(wù)隊(duì)列的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • 利用JavaMail發(fā)送HTML模板郵件

    利用JavaMail發(fā)送HTML模板郵件

    這篇文章主要為大家詳細(xì)介紹了利用JavaMail發(fā)送HTML模板郵件,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2022-08-08
  • SpringBoot構(gòu)建ORM框架的方法步驟

    SpringBoot構(gòu)建ORM框架的方法步驟

    本文主要介紹了SpringBoot構(gòu)建ORM框架的方法步驟,文中通過示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2022-02-02
  • Java定時(shí)調(diào)用.ktr文件的示例代碼(解決方案)

    Java定時(shí)調(diào)用.ktr文件的示例代碼(解決方案)

    這篇文章主要介紹了Java定時(shí)調(diào)用.ktr文件的示例代碼,本文給大家分享遇到問題及解決方法,對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-04-04
  • SpringCache之 @CachePut的使用

    SpringCache之 @CachePut的使用

    這篇文章主要介紹了SpringCache之 @CachePut的使用,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過來看看吧
    2021-02-02
  • linux下idea、pycharm等輸入中文拼音時(shí)滿3個(gè)字母后無法繼續(xù)拼音輸入的問題

    linux下idea、pycharm等輸入中文拼音時(shí)滿3個(gè)字母后無法繼續(xù)拼音輸入的問題

    這篇文章主要介紹了linux下idea、pycharm等輸入中文拼音時(shí)滿3個(gè)字母后無法繼續(xù)拼音輸入的問題,本文通過圖文并茂的形式給大家分享解決方法,需要的朋友可以參考下
    2021-04-04
  • springcloud 熔斷監(jiān)控Hystrix Dashboard和Turbine

    springcloud 熔斷監(jiān)控Hystrix Dashboard和Turbine

    這篇文章主要介紹了springcloud 熔斷監(jiān)控Hystrix Dashboard和Turbine,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-08-08
  • java并發(fā)問題概述

    java并發(fā)問題概述

    這篇文章主要介紹了java并發(fā)問題概述,具有一定借鑒價(jià)值,需要的朋友可以參考下。
    2017-12-12
  • Java后臺(tái)返回blob格式的文件流的解決方案

    Java后臺(tái)返回blob格式的文件流的解決方案

    在Java后臺(tái)開發(fā)中,經(jīng)常會(huì)遇到需要返回Blob格式的文件流給前端的情況,Blob是一種二進(jìn)制大對(duì)象類型,可以用于存儲(chǔ)大量的二進(jìn)制數(shù)據(jù),例如圖片、音頻、視頻等,本文將為你詳細(xì)介紹如何在Java后臺(tái)中返回Blob格式的文件流,需要的朋友可以參考下
    2024-08-08
  • javax.net.ssl.SSLException: java.lang.RuntimeException: Could not generate DH keypair 解決方法總結(jié)

    javax.net.ssl.SSLException: java.lang.RuntimeException: Coul

    這篇文章主要介紹了javax.net.ssl.SSLException: java.lang.RuntimeException: Could not generate DH keypair 解決方法,有需要的朋友們可以學(xué)習(xí)下。
    2019-08-08
  • String類下compareTo()與compare()方法比較

    String類下compareTo()與compare()方法比較

    這篇文章主要介紹了String類下compareTo()與compare()方法比較的相關(guān)資料,需要的朋友可以參考下
    2017-05-05

最新評(píng)論

乐亭县| 怀仁县| 丰顺县| 莫力| 盈江县| 芦山县| 大安市| 灵台县| 灵台县| 开原市| 象州县| 昔阳县| 敦化市| 略阳县| 辛集市| 屏东县| 台湾省| 静海县| 乌兰浩特市| 新营市| 松潘县| 梓潼县| 噶尔县| 蒙自县| 咸阳市| 格尔木市| 乌审旗| 宁河县| 南郑县| 庆阳市| 株洲市| 加查县| 肥城市| 龙江县| 简阳市| 松溪县| 军事| 杭州市| 红原县| 海盐县| 玉环县|