From 57fe3ab0ae9227a4ade102cec532e55f19c6b5d0 Mon Sep 17 00:00:00 2001 From: jaysunxiao Date: Sat, 26 Jun 2021 13:02:00 +0800 Subject: [PATCH] =?UTF-8?q?perf[scheduler]:=20=E4=BC=98=E5=8C=96=E4=BB=A3?= =?UTF-8?q?=E7=A0=81=E7=BB=93=E6=9E=84=EF=BC=8C=E5=B9=B6=E5=88=A0=E9=99=A4?= =?UTF-8?q?=E6=97=A0=E7=94=A8=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../event/schema/EventRegisterProcessor.java | 4 +- .../consumer/registry/ZookeeperRegistry.java | 6 +- .../net/dispatcher/manager/PacketBus.java | 17 ++++- .../java/com/zfoo/net/util/SimpleCache.java | 4 +- .../orm/model/persister/CronOrmPersister.java | 4 +- .../orm/model/persister/TimeOrmPersister.java | 4 +- .../com/zfoo/scheduler/SchedulerContext.java | 17 +---- .../scheduler/manager/ISchedulerManager.java | 40 ---------- ...chedulerManager.java => SchedulerBus.java} | 75 +++++-------------- .../manager/SchedulerThreadFactory.java | 13 ++++ .../schema/SchedulerDefinitionParser.java | 2 - .../schema/SchedulerRegisterProcessor.java | 55 ++++++++++++-- .../com/zfoo/scheduler/util/TimeUtils.java | 6 +- 13 files changed, 114 insertions(+), 133 deletions(-) delete mode 100644 scheduler/src/main/java/com/zfoo/scheduler/manager/ISchedulerManager.java rename scheduler/src/main/java/com/zfoo/scheduler/manager/{SchedulerManager.java => SchedulerBus.java} (68%) diff --git a/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java b/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java index 18501701..a34d00af 100644 --- a/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java +++ b/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java @@ -80,8 +80,8 @@ public class EventRegisterProcessor implements BeanPostProcessor { var enhanceReceiverDefinition = EnhanceUtils.createEventReceiver(receiverDefinition); EventBus.registerEventReceiver(eventClazz, enhanceReceiverDefinition); } - } catch (Exception e) { - throw new RuntimeException(e); + } catch (Throwable t) { + throw new RuntimeException(t); } return bean; diff --git a/net/src/main/java/com/zfoo/net/consumer/registry/ZookeeperRegistry.java b/net/src/main/java/com/zfoo/net/consumer/registry/ZookeeperRegistry.java index 3c3e4802..8a8d26ee 100644 --- a/net/src/main/java/com/zfoo/net/consumer/registry/ZookeeperRegistry.java +++ b/net/src/main/java/com/zfoo/net/consumer/registry/ZookeeperRegistry.java @@ -27,7 +27,7 @@ import com.zfoo.protocol.util.AssertionUtils; import com.zfoo.protocol.util.IOUtils; import com.zfoo.protocol.util.JsonUtils; import com.zfoo.protocol.util.StringUtils; -import com.zfoo.scheduler.SchedulerContext; +import com.zfoo.scheduler.manager.SchedulerBus; import com.zfoo.util.ThreadUtils; import com.zfoo.util.net.HostAndPort; import io.netty.util.concurrent.FastThreadLocalThread; @@ -317,7 +317,7 @@ public class ZookeeperRegistry implements IRegistry { initConsumerCache(); } catch (Exception e) { logger.error("zookeeper初始化失败,等待[{}]秒,重新初始化", RETRY_SECONDS, e); - SchedulerContext.getSchedulerManager().schedule(new Runnable() { + SchedulerBus.schedule(new Runnable() { @Override public void run() { initZookeeper(); @@ -441,7 +441,7 @@ public class ZookeeperRegistry implements IRegistry { } if (recheckFlag) { - SchedulerContext.getSchedulerManager().schedule(new Runnable() { + SchedulerBus.schedule(new Runnable() { @Override public void run() { checkConsumer(); diff --git a/net/src/main/java/com/zfoo/net/dispatcher/manager/PacketBus.java b/net/src/main/java/com/zfoo/net/dispatcher/manager/PacketBus.java index aceceb3c..7a8c2f1c 100644 --- a/net/src/main/java/com/zfoo/net/dispatcher/manager/PacketBus.java +++ b/net/src/main/java/com/zfoo/net/dispatcher/manager/PacketBus.java @@ -1,3 +1,16 @@ +/* + * Copyright (C) 2020 The zfoo Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except + * in compliance with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed + * on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and limitations under the License. + */ + package com.zfoo.net.dispatcher.manager; import com.zfoo.event.model.event.IEvent; @@ -111,8 +124,8 @@ public abstract class PacketBus { var receiverDefinition = new PacketReceiverDefinition(bean, method, packetClazz, attachmentClazz); var enhanceReceiverDefinition = EnhanceUtils.createPacketReceiver(receiverDefinition); packetReceiverList[protocolId] = enhanceReceiverDefinition; - } catch (Exception e) { - throw new RuntimeException(e); + } catch (Throwable t) { + throw new RuntimeException(t); } } } diff --git a/net/src/main/java/com/zfoo/net/util/SimpleCache.java b/net/src/main/java/com/zfoo/net/util/SimpleCache.java index f50912ca..2c10c522 100644 --- a/net/src/main/java/com/zfoo/net/util/SimpleCache.java +++ b/net/src/main/java/com/zfoo/net/util/SimpleCache.java @@ -19,7 +19,7 @@ import com.github.benmanes.caffeine.cache.LoadingCache; import com.zfoo.event.manager.EventBus; import com.zfoo.protocol.collection.CollectionUtils; import com.zfoo.protocol.model.Pair; -import com.zfoo.scheduler.manager.SchedulerManager; +import com.zfoo.scheduler.manager.SchedulerBus; import org.checkerframework.checker.nullness.qual.NonNull; import org.checkerframework.checker.nullness.qual.Nullable; @@ -92,7 +92,7 @@ public class SimpleCache { }); - SchedulerManager.getInstance().scheduleAtFixedRate(new Runnable() { + SchedulerBus.scheduleAtFixedRate(new Runnable() { @Override public void run() { // 不在任务调度线程中执行耗时任务,因为任务调度线程只有一个线程池 diff --git a/orm/src/main/java/com/zfoo/orm/model/persister/CronOrmPersister.java b/orm/src/main/java/com/zfoo/orm/model/persister/CronOrmPersister.java index 186c45f6..47bb4078 100644 --- a/orm/src/main/java/com/zfoo/orm/model/persister/CronOrmPersister.java +++ b/orm/src/main/java/com/zfoo/orm/model/persister/CronOrmPersister.java @@ -18,7 +18,7 @@ import com.zfoo.orm.OrmContext; import com.zfoo.orm.model.cache.EntityCaches; import com.zfoo.orm.model.vo.EntityDef; import com.zfoo.protocol.exception.ExceptionUtils; -import com.zfoo.scheduler.SchedulerContext; +import com.zfoo.scheduler.manager.SchedulerBus; import com.zfoo.scheduler.util.TimeUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -74,7 +74,7 @@ public class CronOrmPersister extends AbstractOrmPersister { } if (!OrmContext.isStop()) { - SchedulerContext.getSchedulerManager().schedule(new Runnable() { + SchedulerBus.schedule(new Runnable() { @Override public void run() { if (!OrmContext.isStop()) { diff --git a/orm/src/main/java/com/zfoo/orm/model/persister/TimeOrmPersister.java b/orm/src/main/java/com/zfoo/orm/model/persister/TimeOrmPersister.java index 6c00d907..4c2493c2 100644 --- a/orm/src/main/java/com/zfoo/orm/model/persister/TimeOrmPersister.java +++ b/orm/src/main/java/com/zfoo/orm/model/persister/TimeOrmPersister.java @@ -18,7 +18,7 @@ import com.zfoo.orm.OrmContext; import com.zfoo.orm.model.cache.EntityCaches; import com.zfoo.orm.model.vo.EntityDef; import com.zfoo.protocol.util.StringUtils; -import com.zfoo.scheduler.SchedulerContext; +import com.zfoo.scheduler.manager.SchedulerBus; import java.util.concurrent.TimeUnit; @@ -43,7 +43,7 @@ public class TimeOrmPersister extends AbstractOrmPersister { @Override public void start() { - SchedulerContext.getSchedulerManager().scheduleAtFixedRate(new Runnable() { + SchedulerBus.scheduleAtFixedRate(new Runnable() { @Override public void run() { if (!OrmContext.isStop()) { diff --git a/scheduler/src/main/java/com/zfoo/scheduler/SchedulerContext.java b/scheduler/src/main/java/com/zfoo/scheduler/SchedulerContext.java index 2306b910..1318d3fa 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/SchedulerContext.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/SchedulerContext.java @@ -14,8 +14,7 @@ package com.zfoo.scheduler; import com.zfoo.protocol.util.ReflectionUtils; -import com.zfoo.scheduler.manager.ISchedulerManager; -import com.zfoo.scheduler.manager.SchedulerManager; +import com.zfoo.scheduler.manager.SchedulerBus; import com.zfoo.scheduler.schema.SchedulerRegisterProcessor; import com.zfoo.util.ThreadUtils; import org.slf4j.Logger; @@ -53,11 +52,6 @@ public class SchedulerContext implements ApplicationListener schedulerDefList = new CopyOnWriteArrayList<>(); @@ -56,10 +50,10 @@ public class SchedulerManager implements ISchedulerManager { private static long minSchedulerTriggerTimestamp = 0; - private static final long TRIGGER_MILLIS_INTERVAL = TimeUtils.MILLIS_PER_SECOND; + public static final long TRIGGER_MILLIS_INTERVAL = TimeUtils.MILLIS_PER_SECOND; /** - * scheduler默认只有一个单线程线程池 + * scheduler默认只有一个单线程的线程池 */ private static final ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(new SchedulerThreadFactory(1)); @@ -71,15 +65,9 @@ public class SchedulerManager implements ISchedulerManager { } catch (Exception e) { logger.error("scheduler triggers an error.", e); } - }, TimeUtils.MILLIS_PER_SECOND, TRIGGER_MILLIS_INTERVAL, TimeUnit.MILLISECONDS); + }, 3 * TimeUtils.MILLIS_PER_SECOND, TRIGGER_MILLIS_INTERVAL, TimeUnit.MILLISECONDS); } - public static SchedulerManager getInstance() { - return INSTANCE; - } - - private SchedulerManager() { - } private static long minSchedulerTriggerTimestamp() { var minSchedulerOptional = schedulerDefList.stream().min(Comparator.comparingLong(schedulerDef -> schedulerDef.getTriggerTimestamp())); @@ -133,44 +121,16 @@ public class SchedulerManager implements ISchedulerManager { minSchedulerTriggerTimestamp = minSchedulerTriggerTimestamp(); } - public void registerScheduler(Object bean) { - try { - var methods = ReflectionUtils.getMethodsByAnnoInPOJOClass(bean.getClass(), Scheduler.class); - for (var method : methods) { - var scheduler = method.getAnnotation(Scheduler.class); - - var paramClazzs = method.getParameterTypes(); - if (paramClazzs.length >= 1) { - throw new IllegalArgumentException(StringUtils.format("[class:{}] [method:{}] can not have any parameters", bean.getClass(), method.getName())); - } - - var methodName = method.getName(); - - if (!Modifier.isPublic(method.getModifiers())) { - throw new IllegalArgumentException(StringUtils.format("[class:{}] [method:{}] must use 'public' as modifier!", bean.getClass().getName(), methodName)); - } - - if (Modifier.isStatic(method.getModifiers())) { - throw new IllegalArgumentException(StringUtils.format("[class:{}] [method:{}] can not use 'static' as modifier!", bean.getClass().getName(), methodName)); - } - - if (!methodName.startsWith("cron")) { - throw new IllegalArgumentException(StringUtils.format("[class:{}] [method:{}] must start with 'cron' as method name!" - , bean.getClass().getName(), methodName)); - } - - var schedulerDef = SchedulerDefinition.valueOf(scheduler.cron(), bean, method); - schedulerDefList.add(schedulerDef); - minSchedulerTriggerTimestamp = minSchedulerTriggerTimestamp(); - } - } catch (Throwable throwable) { - throw new RuntimeException(throwable); - } + public static void registerScheduler(SchedulerDefinition scheduler) { + schedulerDefList.add(scheduler); + minSchedulerTriggerTimestamp = minSchedulerTriggerTimestamp(); } - @Override - public void scheduleAtFixedRate(Runnable runnable, long period, TimeUnit unit) { + /** + * 不断执行的周期循环任务 + */ + public static void scheduleAtFixedRate(Runnable runnable, long period, TimeUnit unit) { if (SchedulerContext.isStop()) { return; } @@ -189,8 +149,11 @@ public class SchedulerManager implements ISchedulerManager { }, 0, period, unit); } - @Override - public void schedule(Runnable runnable, long delay, TimeUnit unit) { + + /** + * 固定延迟执行的任务 + */ + public static void schedule(Runnable runnable, long delay, TimeUnit unit) { if (SchedulerContext.isStop()) { return; } @@ -209,7 +172,9 @@ public class SchedulerManager implements ISchedulerManager { }, delay, unit); } - @Override + /** + * cron表达式执行的任务 + */ public void scheduleCron(Runnable runnable, String cron) { if (SchedulerContext.isStop()) { return; diff --git a/scheduler/src/main/java/com/zfoo/scheduler/manager/SchedulerThreadFactory.java b/scheduler/src/main/java/com/zfoo/scheduler/manager/SchedulerThreadFactory.java index 30b3ccc9..8a646b22 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/manager/SchedulerThreadFactory.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/manager/SchedulerThreadFactory.java @@ -1,3 +1,16 @@ +/* + * Copyright (C) 2020 The zfoo Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except + * in compliance with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed + * on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and limitations under the License. + */ + package com.zfoo.scheduler.manager; import io.netty.util.concurrent.FastThreadLocalThread; diff --git a/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerDefinitionParser.java b/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerDefinitionParser.java index 169ffc94..52749559 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerDefinitionParser.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerDefinitionParser.java @@ -47,8 +47,6 @@ public class SchedulerDefinitionParser implements BeanDefinitionParser { builder = BeanDefinitionBuilder.rootBeanDefinition(clazz); parserContext.getRegistry().registerBeanDefinition(name, builder.getBeanDefinition()); - // 注册SchedulerManager - String schedulerId = element.getAttribute(SCHEDULER_ID); return builder.getBeanDefinition(); } diff --git a/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerRegisterProcessor.java b/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerRegisterProcessor.java index 1e594e9d..d8c91776 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerRegisterProcessor.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerRegisterProcessor.java @@ -13,25 +13,70 @@ package com.zfoo.scheduler.schema; -import com.zfoo.scheduler.SchedulerContext; -import com.zfoo.scheduler.manager.SchedulerManager; +import com.zfoo.protocol.collection.ArrayUtils; +import com.zfoo.protocol.util.ReflectionUtils; +import com.zfoo.protocol.util.StringUtils; +import com.zfoo.scheduler.manager.SchedulerBus; +import com.zfoo.scheduler.model.anno.Scheduler; +import com.zfoo.scheduler.model.vo.SchedulerDefinition; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.beans.BeansException; import org.springframework.beans.factory.config.BeanPostProcessor; +import java.lang.reflect.Modifier; + /** * @author jaysunxiao * @version 3.0 */ public class SchedulerRegisterProcessor implements BeanPostProcessor { + private static final Logger logger = LoggerFactory.getLogger(SchedulerRegisterProcessor.class); @Override public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { - if (SchedulerContext.getSchedulerContext() == null) { + var clazz = bean.getClass(); + var methods = ReflectionUtils.getMethodsByAnnoInPOJOClass(bean.getClass(), Scheduler.class); + + if (ArrayUtils.isEmpty(methods)) { return bean; } - SchedulerManager schedulerManager = (SchedulerManager) SchedulerContext.getSchedulerManager(); - schedulerManager.registerScheduler(bean); + + if (!ReflectionUtils.isPojoClass(clazz)) { + logger.warn("调度注册类[{}]不是POJO类,父类的调度不会被扫描到", clazz); + } + + try { + for (var method : methods) { + var schedulerMethod = method.getAnnotation(Scheduler.class); + + var paramClazzs = method.getParameterTypes(); + if (paramClazzs.length >= 1) { + throw new IllegalArgumentException(StringUtils.format("[class:{}] [method:{}] can not have any parameters", bean.getClass(), method.getName())); + } + + var methodName = method.getName(); + + if (!Modifier.isPublic(method.getModifiers())) { + throw new IllegalArgumentException(StringUtils.format("[class:{}] [method:{}] must use 'public' as modifier!", bean.getClass().getName(), methodName)); + } + + if (Modifier.isStatic(method.getModifiers())) { + throw new IllegalArgumentException(StringUtils.format("[class:{}] [method:{}] can not use 'static' as modifier!", bean.getClass().getName(), methodName)); + } + + if (!methodName.startsWith("cron")) { + throw new IllegalArgumentException(StringUtils.format("[class:{}] [method:{}] must start with 'cron' as method name!" + , bean.getClass().getName(), methodName)); + } + + var scheduler = SchedulerDefinition.valueOf(schedulerMethod.cron(), bean, method); + SchedulerBus.registerScheduler(scheduler); + } + } catch (Throwable t) { + throw new RuntimeException(t); + } return bean; } diff --git a/scheduler/src/main/java/com/zfoo/scheduler/util/TimeUtils.java b/scheduler/src/main/java/com/zfoo/scheduler/util/TimeUtils.java index 44a06d04..7d335d97 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/util/TimeUtils.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/util/TimeUtils.java @@ -13,8 +13,7 @@ package com.zfoo.scheduler.util; -import com.zfoo.scheduler.SchedulerContext; -import com.zfoo.scheduler.manager.SchedulerManager; +import com.zfoo.scheduler.manager.SchedulerBus; import io.netty.util.concurrent.FastThreadLocal; import org.springframework.scheduling.support.CronExpression; @@ -81,8 +80,7 @@ public abstract class TimeUtils { }; static { - SchedulerContext.getSchedulerManager(); - SchedulerManager.getInstance(); + var interval = SchedulerBus.TRIGGER_MILLIS_INTERVAL; } /**