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

DUCC配置平臺實(shí)現(xiàn)一個(gè)動(dòng)態(tài)化線程池示例代碼

 更新時(shí)間:2023年02月16日 16:54:40   作者:京東云開發(fā)者  
這篇文章主要為大家介紹了DUCC配置平臺實(shí)現(xiàn)一個(gè)動(dòng)態(tài)化線程池示例代碼,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

作者:京東零售 張賓

1.背景

在后臺開發(fā)中,會經(jīng)常用到線程池技術(shù),對于線程池核心參數(shù)的配置很大程度上依靠經(jīng)驗(yàn)。然而,由于系統(tǒng)運(yùn)行過程中存在的不確定性,我們很難一勞永逸地規(guī)劃一個(gè)合理的線程池參數(shù)。在對線程池配置參數(shù)進(jìn)行調(diào)整時(shí),一般需要對服務(wù)進(jìn)行重啟,這樣修改的成本就會偏高。一種解決辦法就是,將線程池的配置放到配置平臺側(cè),系統(tǒng)運(yùn)行期間開發(fā)人員根據(jù)系統(tǒng)運(yùn)行情況對核心參數(shù)進(jìn)行動(dòng)態(tài)配置。

本文以公司DUCC配置平臺作為服務(wù)配置中心,以修改線程池核心線程數(shù)、最大線程數(shù)為例,實(shí)現(xiàn)一個(gè)簡單的動(dòng)態(tài)化線程池。

2.代碼實(shí)現(xiàn)

當(dāng)前項(xiàng)目中使用的是Spring 框架提供的線程池類ThreadPoolTaskExecutor,而ThreadPoolTaskExecutor底層又使用里了JDK中線程池類ThreadPoolExecutor,線程池類ThreadPoolExecutor有兩個(gè)成員方法setCorePoolSize、setMaximumPoolSize可以在運(yùn)行時(shí)設(shè)置核心線程數(shù)和最大線程數(shù)。

setCorePoolSize方法執(zhí)行流程是:首先會覆蓋之前構(gòu)造函數(shù)設(shè)置的corePoolSize,然后,如果新的值比原始值要小,當(dāng)多余的工作線程下次變成空閑狀態(tài)的時(shí)候會被中斷并銷毀,如果新的值比原來的值要大且工作隊(duì)列不為空,則會創(chuàng)建新的工作線程。流程圖如下:

setMaximumPoolSize方法: 首先會覆蓋之前構(gòu)造函數(shù)設(shè)置的maximumPoolSize,然后,如果新的值比原來的值要小,當(dāng)多余的工作線程下次變成空閑狀態(tài)的時(shí)候會被中斷并銷毀。

Spring 框架提供的線程池類ThreadPoolTaskExecutor,此類封裝了對ThreadPoolExecutor有兩個(gè)成員方法setCorePoolSize、setMaximumPoolSize的調(diào)用。

基于以上源代碼分析,要實(shí)現(xiàn)一個(gè)簡單的動(dòng)態(tài)線程池需要以下幾步:

(1)定義一個(gè)動(dòng)態(tài)線程池類,繼承ThreadPoolTaskExecutor,目的跟非動(dòng)態(tài)配置的線程池類ThreadPoolTaskExecutor區(qū)分開;

(2)定義和實(shí)現(xiàn)一個(gè)動(dòng)態(tài)線程池配置定時(shí)刷的類,目的定時(shí)對比ducc配置的線程池?cái)?shù)和本地應(yīng)用中線程數(shù)是否一致,若不一致,則更新本地動(dòng)態(tài)線程池線程池?cái)?shù);

(3)引入公司ducc配置平臺相關(guān)jar包并創(chuàng)建一個(gè)動(dòng)態(tài)線程池配置key;

(4)定義和實(shí)現(xiàn)一個(gè)應(yīng)用啟動(dòng)后根據(jù)動(dòng)態(tài)線程池Bean和從ducc配置平臺拉取配置刷新應(yīng)用中的線程數(shù)配置;

接下來代碼一一實(shí)現(xiàn):

(1)動(dòng)態(tài)線程池類

/**
 * 動(dòng)態(tài)線程池
 *
 */
public class DynamicThreadPoolTaskExecutor extends ThreadPoolTaskExecutor {
}

(2)動(dòng)態(tài)線程池配置定時(shí)刷新類

@Slf4j
public class DynamicThreadPoolRefresh implements InitializingBean {
    /**
     * Maintain all automatically registered and manually registered DynamicThreadPoolTaskExecutor.
     */
    private static final ConcurrentMap<String, DynamicThreadPoolTaskExecutor> DTP_REGISTRY = new ConcurrentHashMap<>();
    /**
     * @param threadPoolBeanName
     * @param threadPoolTaskExecutor
     */
    public static void registerDynamicThreadPool(String threadPoolBeanName, DynamicThreadPoolTaskExecutor threadPoolTaskExecutor) {
        log.info("DynamicThreadPool register ThreadPoolTaskExecutor, threadPoolBeanName: {}, executor: {}", threadPoolBeanName, ExecutorConverter.convert(threadPoolBeanName, threadPoolTaskExecutor.getThreadPoolExecutor()));
        DTP_REGISTRY.putIfAbsent(threadPoolBeanName, threadPoolTaskExecutor);
    }
    @Override
    public void afterPropertiesSet() throws Exception {
        this.refresh();
        //創(chuàng)建定時(shí)任務(wù)線程池
        ScheduledExecutorService executorService = new ScheduledThreadPoolExecutor(1, (new BasicThreadFactory.Builder()).namingPattern("DynamicThreadPoolRefresh-%d").daemon(true).build());
        //延遲1秒執(zhí)行,每個(gè)1分鐘check一次
        executorService.scheduleAtFixedRate(new RefreshThreadPoolConfig(), 1000L, 60000L, TimeUnit.MILLISECONDS);
    }
    private void refresh() {
        String dynamicThreadPool = "";
        try {
            if (DTP_REGISTRY.isEmpty()) {
                log.debug("DynamicThreadPool refresh DTP_REGISTRY is empty");
                return;
            }
            dynamicThreadPool = DuccConfigUtil.getValue(DuccConfigConstants.DYNAMIC_THREAD_POOL);
            if (StringUtils.isBlank(dynamicThreadPool)) {
                log.debug("DynamicThreadPool refresh dynamicThreadPool not config");
                return;
            }
            log.debug("DynamicThreadPool refresh dynamicThreadPool:{}", dynamicThreadPool);
            List<ThreadPoolProperties> threadPoolPropertiesList = JsonUtil.json2Object(dynamicThreadPool, new TypeReference<List<ThreadPoolProperties>>() {
            });
            if (CollectionUtils.isEmpty(threadPoolPropertiesList)) {
                log.error("DynamicThreadPool refresh dynamicThreadPool json2Object error!{}", dynamicThreadPool);
                return;
            }
            for (ThreadPoolProperties properties : threadPoolPropertiesList) {
                doRefresh(properties);
            }
        } catch (Exception e) {
            log.error("DynamicThreadPool refresh exception!dynamicThreadPool:{}", dynamicThreadPool, e);
        }
    }
    /**
     * @param properties
     */
    private void doRefresh(ThreadPoolProperties properties) {
        if (StringUtils.isBlank(properties.getThreadPoolBeanName())
                || properties.getCorePoolSize() < 1
                || properties.getMaxPoolSize() < 1
                || properties.getMaxPoolSize() < properties.getCorePoolSize()) {
            log.error("DynamicThreadPool refresh, invalid parameters exist, properties: {}", properties);
            return;
        }
        DynamicThreadPoolTaskExecutor threadPoolTaskExecutor = DTP_REGISTRY.get(properties.getThreadPoolBeanName());
        if (Objects.isNull(threadPoolTaskExecutor)) {
            log.warn("DynamicThreadPool refresh, DTP_REGISTRY not found {}", properties.getThreadPoolBeanName());
            return;
        }
        ThreadPoolProperties oldProp = ExecutorConverter.convert(properties.getThreadPoolBeanName(), threadPoolTaskExecutor.getThreadPoolExecutor());
        if (Objects.equals(oldProp.getCorePoolSize(), properties.getCorePoolSize())
                && Objects.equals(oldProp.getMaxPoolSize(), properties.getMaxPoolSize())) {
            log.warn("DynamicThreadPool refresh, properties of [{}] have not changed.", properties.getThreadPoolBeanName());
            return;
        }
        if (!Objects.equals(oldProp.getCorePoolSize(), properties.getCorePoolSize())) {
            threadPoolTaskExecutor.setCorePoolSize(properties.getCorePoolSize());
            log.info("DynamicThreadPool refresh, corePoolSize changed!{} {}", properties.getThreadPoolBeanName(), properties.getCorePoolSize());
        }
        if (!Objects.equals(oldProp.getMaxPoolSize(), properties.getMaxPoolSize())) {
            threadPoolTaskExecutor.setMaxPoolSize(properties.getMaxPoolSize());
            log.info("DynamicThreadPool refresh, maxPoolSize changed!{} {}", properties.getThreadPoolBeanName(), properties.getMaxPoolSize());
        }
        ThreadPoolProperties newProp = ExecutorConverter.convert(properties.getThreadPoolBeanName(), threadPoolTaskExecutor.getThreadPoolExecutor());
        log.info("DynamicThreadPool refresh result!{} oldProp:{},newProp:{}", properties.getThreadPoolBeanName(), oldProp, newProp);
    }
    private class RefreshThreadPoolConfig extends TimerTask {
        private RefreshThreadPoolConfig() {
        }
        @Override
        public void run() {
            DynamicThreadPoolRefresh.this.refresh();
        }
    }
}

線程池配置類

@Data
public class ThreadPoolProperties {
    /**
     * 線程池名稱
     */
    private String threadPoolBeanName;
    /**
     * 線程池核心線程數(shù)量
     */
    private int corePoolSize;
    /**
     * 線程池最大線程池?cái)?shù)量
     */
    private int maxPoolSize;
}

(3)引入公司ducc配置平臺相關(guān)jar包并創(chuàng)建一個(gè)動(dòng)態(tài)線程池配置key

ducc配置平臺使用見:cf.jd.com/pages/viewp…?

動(dòng)態(tài)線程池配置key:dynamic.thread.pool

配置value:

[  {    "threadPoolBeanName": "submitOrderThreadPoolTaskExecutor",    "corePoolSize": 32,    "maxPoolSize": 128  }]

(4) 應(yīng)用啟動(dòng)刷新應(yīng)用本地動(dòng)態(tài)線程池配置

@Slf4j
public class DynamicThreadPoolPostProcessor implements BeanPostProcessor {
    @Override
    public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
        if (bean instanceof DynamicThreadPoolTaskExecutor) {
            DynamicThreadPoolRefresh.registerDynamicThreadPool(beanName, (DynamicThreadPoolTaskExecutor) bean);
        }
        return bean;
    }
}

3.動(dòng)態(tài)線程池應(yīng)用

動(dòng)態(tài)線程池Bean聲明

    <!-- 普通線程池 -->
    <bean id="threadPoolTaskExecutor" class="com.jd.concurrent.ThreadPoolTaskExecutorWrapper">
        <!-- 核心線程數(shù),默認(rèn)為 -->
        <property name="corePoolSize" value="128"/>
        <!-- 最大線程數(shù),默認(rèn)為Integer.MAX_VALUE -->
        <property name="maxPoolSize" value="512"/>
        <!-- 隊(duì)列最大長度,一般需要設(shè)置值>=notifyScheduledMainExecutor.maxNum;默認(rèn)為Integer.MAX_VALUE -->
        <property name="queueCapacity" value="500"/>
        <!-- 線程池維護(hù)線程所允許的空閑時(shí)間,默認(rèn)為60s -->
        <property name="keepAliveSeconds" value="60"/>
        <!-- 線程池對拒絕任務(wù)(無線程可用)的處理策略,目前只支持AbortPolicy、CallerRunsPolicy;默認(rèn)為后者 -->
        <property name="rejectedExecutionHandler">
            <!-- AbortPolicy:直接拋出java.util.concurrent.RejectedExecutionException異常 -->
            <!-- CallerRunsPolicy:主線程直接執(zhí)行該任務(wù),執(zhí)行完之后嘗試添加下一個(gè)任務(wù)到線程池中,可以有效降低向線程池內(nèi)添加任務(wù)的速度 -->
            <!-- DiscardOldestPolicy:拋棄舊的任務(wù)、暫不支持;會導(dǎo)致被丟棄的任務(wù)無法再次被執(zhí)行 -->
            <!-- DiscardPolicy:拋棄當(dāng)前任務(wù)、暫不支持;會導(dǎo)致被丟棄的任務(wù)無法再次被執(zhí)行 -->
            <bean class="java.util.concurrent.ThreadPoolExecutor$CallerRunsPolicy"/>
        </property>
    </bean>
    <!-- 動(dòng)態(tài)線程池 -->
    <bean id="submitOrderThreadPoolTaskExecutor" class="com.jd.concurrent.DynamicThreadPoolTaskExecutor">
        <!-- 核心線程數(shù),默認(rèn)為 -->
        <property name="corePoolSize" value="32"/>
        <!-- 最大線程數(shù),默認(rèn)為Integer.MAX_VALUE -->
        <property name="maxPoolSize" value="128"/>
        <!-- 隊(duì)列最大長度,一般需要設(shè)置值>=notifyScheduledMainExecutor.maxNum;默認(rèn)為Integer.MAX_VALUE -->
        <property name="queueCapacity" value="500"/>
        <!-- 線程池維護(hù)線程所允許的空閑時(shí)間,默認(rèn)為60s -->
        <property name="keepAliveSeconds" value="60"/>
        <!-- 線程池對拒絕任務(wù)(無線程可用)的處理策略,目前只支持AbortPolicy、CallerRunsPolicy;默認(rèn)為后者 -->
        <property name="rejectedExecutionHandler">
            <!-- AbortPolicy:直接拋出java.util.concurrent.RejectedExecutionException異常 -->
            <!-- CallerRunsPolicy:主線程直接執(zhí)行該任務(wù),執(zhí)行完之后嘗試添加下一個(gè)任務(wù)到線程池中,可以有效降低向線程池內(nèi)添加任務(wù)的速度 -->
            <!-- DiscardOldestPolicy:拋棄舊的任務(wù)、暫不支持;會導(dǎo)致被丟棄的任務(wù)無法再次被執(zhí)行 -->
            <!-- DiscardPolicy:拋棄當(dāng)前任務(wù)、暫不支持;會導(dǎo)致被丟棄的任務(wù)無法再次被執(zhí)行 -->
            <bean class="java.util.concurrent.ThreadPoolExecutor$CallerRunsPolicy"/>
        </property>
    </bean>
    <!-- 動(dòng)態(tài)線程池刷新配置 -->
    <bean class="com.jd.concurrent.DynamicThreadPoolPostProcessor"/>
    <bean class="com.jd.concurrent.DynamicThreadPoolRefresh"/>

業(yè)務(wù)類注入Spring Bean后,直接使用即可

 @Resource
 private ThreadPoolTaskExecutor submitOrderThreadPoolTaskExecutor;
 Runnable asyncTask = ()-&gt;{...};
 CompletableFuture.runAsync(asyncTask, this.submitOrderThreadPoolTaskExecutor);

4.小結(jié)

本文從實(shí)際項(xiàng)目的業(yè)務(wù)痛點(diǎn)場景出發(fā),并基于公司已有的ducc配置平臺簡單實(shí)現(xiàn)了線程池線程數(shù)量可配置。

以上就是DUCC配置平臺實(shí)現(xiàn)一個(gè)動(dòng)態(tài)化線程池示例代碼的詳細(xì)內(nèi)容,更多關(guān)于DUCC配置平臺動(dòng)態(tài)化線程池的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • SpringMVC適配器模式代碼示例

    SpringMVC適配器模式代碼示例

    這篇文章主要介紹了SpringMVC適配器模式代碼示例,涉及模擬springmvc的Java代碼等相關(guān)內(nèi)容,具有一定借鑒價(jià)值,需要的朋友可以參考下。
    2017-11-11
  • Java算法實(shí)現(xiàn)楊輝三角的講解

    Java算法實(shí)現(xiàn)楊輝三角的講解

    今天小編就為大家分享一篇關(guān)于Java算法實(shí)現(xiàn)楊輝三角的講解,小編覺得內(nèi)容挺不錯(cuò)的,現(xiàn)在分享給大家,具有很好的參考價(jià)值,需要的朋友一起跟隨小編來看看吧
    2019-01-01
  • 使用JPA自定義id策略避免主鍵自增

    使用JPA自定義id策略避免主鍵自增

    這篇文章主要介紹了使用JPA自定義id策略避免主鍵自增問題,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-08-08
  • Java中切面的使用方法舉例詳解

    Java中切面的使用方法舉例詳解

    這篇文章主要介紹了Java中切面編程(AOP)的基本概念、原理及實(shí)現(xiàn)方式,AOP通過將橫切關(guān)注點(diǎn)模塊化為切面,使代碼更易于維護(hù)和擴(kuò)展,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2025-03-03
  • SpringBoot優(yōu)雅地實(shí)現(xiàn)全局異常處理的方法詳解

    SpringBoot優(yōu)雅地實(shí)現(xiàn)全局異常處理的方法詳解

    這篇文章主要為大家詳細(xì)介紹了SpringBoot如何優(yōu)雅地實(shí)現(xiàn)全局異常處理,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下
    2022-08-08
  • Java中的jinfo命令使用詳解

    Java中的jinfo命令使用詳解

    jinfo是JDK提供的一個(gè)可以實(shí)時(shí)查看Java虛擬機(jī)各種配置參數(shù)和系統(tǒng)屬性的命令行工具,本文給大家介紹下Java中的jinfo命令使用,感興趣的朋友一起看看吧
    2022-03-03
  • 使用@Value為靜態(tài)變量導(dǎo)入并使用導(dǎo)入的靜態(tài)變量進(jìn)行初始化方式

    使用@Value為靜態(tài)變量導(dǎo)入并使用導(dǎo)入的靜態(tài)變量進(jìn)行初始化方式

    這篇文章主要介紹了使用@Value為靜態(tài)變量導(dǎo)入并使用導(dǎo)入的靜態(tài)變量進(jìn)行初始化方式,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-02-02
  • java?使用BeanFactory實(shí)現(xiàn)service與dao層解耦合詳解

    java?使用BeanFactory實(shí)現(xiàn)service與dao層解耦合詳解

    這篇文章主要介紹了java?使用BeanFactory實(shí)現(xiàn)service與dao層解耦合詳解,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-12-12
  • Java實(shí)現(xiàn)哈希表的基本功能

    Java實(shí)現(xiàn)哈希表的基本功能

    今天教大家怎么用Java實(shí)現(xiàn)哈希表的基本功能,文中有非常詳細(xì)的代碼示例,對正在學(xué)習(xí)java的小伙伴們有非常好的幫助,需要的朋友可以參考下
    2021-05-05
  • WebService教程詳解(一)

    WebService教程詳解(一)

    WebService,顧名思義就是基于Web的服務(wù)。它使用Web(HTTP)方式,接收和響應(yīng)外部系統(tǒng)的某種請求,接下來通過本文給大家介紹WebService教程詳解(一),對webservice教程感興趣的朋友一起學(xué)習(xí)吧
    2016-03-03

最新評論

无棣县| 阳谷县| 行唐县| 开封县| 郎溪县| 从化市| 长泰县| 社旗县| 长宁区| 岳西县| 星座| 安图县| 通城县| 隆安县| 云南省| 乐亭县| 资阳市| 阿城市| 黔南| 安远县| 开封市| 乳源| 金湖县| 瑞丽市| 镇雄县| 佛教| 滕州市| 乌拉特中旗| 乐清市| 乌什县| 同江市| 西平县| 岱山县| 开阳县| 呈贡县| 兴宁市| 双柏县| 霍山县| 武功县| 澜沧| 浙江省|