From c4cdf272ebf27eaee18f6096434002223827d05b Mon Sep 17 00:00:00 2001 From: godotg Date: Thu, 28 Jul 2022 21:47:33 +0800 Subject: [PATCH] =?UTF-8?q?perf[runnable]:=20=E4=BD=BF=E7=94=A8Runnable?= =?UTF-8?q?=E5=87=BD=E6=95=B0=E5=BC=8F=E7=BC=96=E7=A8=8B=E6=8E=A5=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../java/com/zfoo/event/EventContext.java | 2 +- .../java/com/zfoo/event/manager/EventBus.java | 12 +++--- .../zfoo/event/model/anno/EventReceiver.java | 2 +- .../event/model/event/AppStartAfterEvent.java | 2 +- .../model/event/AppStartBeforeEvent.java | 2 +- .../zfoo/event/model/event/AppStartEvent.java | 2 +- .../com/zfoo/event/model/event/IEvent.java | 2 +- .../com/zfoo/event/model/vo/EnhanceUtils.java | 2 +- .../model/vo/EventReceiverDefinition.java | 2 +- .../zfoo/event/model/vo/IEventReceiver.java | 2 +- .../event/schema/EventDefinitionParser.java | 2 +- .../event/schema/EventRegisterProcessor.java | 2 +- .../zfoo/event/schema/NamespaceHandler.java | 2 +- .../java/com/zfoo/event/ApplicationTest.java | 2 +- .../java/com/zfoo/event/MyController1.java | 2 +- .../java/com/zfoo/event/MyController2.java | 2 +- .../java/com/zfoo/event/MyNoticeEvent.java | 2 +- .../consumer/registry/ZookeeperRegistry.java | 15 +------ .../main/java/com/zfoo/net/task/TaskBus.java | 10 ++--- .../java/com/zfoo/net/util/SimpleCache.java | 43 ++++++++----------- .../java/com/zfoo/net/util/SingleCache.java | 8 +--- .../com/zfoo/net/router/SignalBridgeTest.java | 9 ++-- .../orm/model/persister/CronOrmPersister.java | 21 +++------ .../orm/model/persister/TimeOrmPersister.java | 17 ++------ .../lpmap/ConcurrentFileChannelMapTest.java | 5 +-- .../zfoo/orm/lpmap/ConcurrentHeapMapTest.java | 18 +++----- .../com/zfoo/scheduler/SchedulerContext.java | 2 +- .../zfoo/scheduler/manager/SchedulerBus.java | 10 ++--- .../com/zfoo/scheduler/model/StopWatch.java | 4 +- .../zfoo/scheduler/model/anno/Scheduler.java | 2 +- .../zfoo/scheduler/model/vo/EnhanceUtils.java | 2 +- .../zfoo/scheduler/model/vo/IScheduler.java | 2 +- .../scheduler/model/vo/ReflectScheduler.java | 2 +- .../scheduler/model/vo/RunnableScheduler.java | 2 +- .../model/vo/SchedulerDefinition.java | 2 +- .../scheduler/schema/NamespaceHandler.java | 2 +- .../schema/SchedulerDefinitionParser.java | 2 +- .../com/zfoo/scheduler/util/TimeUtils.java | 2 +- .../com/zfoo/scheduler/ApplicationTest.java | 2 +- .../zfoo/scheduler/SchedulerController.java | 2 +- .../zfoo/scheduler/util/TimeUtilsTest.java | 2 +- .../main/java/com/zfoo/util/SafeRunnable.java | 17 ++++++-- 42 files changed, 105 insertions(+), 142 deletions(-) diff --git a/event/src/main/java/com/zfoo/event/EventContext.java b/event/src/main/java/com/zfoo/event/EventContext.java index 2dbe7d58..5e4c6c9c 100644 --- a/event/src/main/java/com/zfoo/event/EventContext.java +++ b/event/src/main/java/com/zfoo/event/EventContext.java @@ -29,7 +29,7 @@ import org.springframework.core.Ordered; import java.util.concurrent.ExecutorService; /** - * @author jaysunxiao + * @author godotg * @version 3.0 *

* 在EventRegisterProcessor 中完成扫描 diff --git a/event/src/main/java/com/zfoo/event/manager/EventBus.java b/event/src/main/java/com/zfoo/event/manager/EventBus.java index 6593921c..cdc3d865 100644 --- a/event/src/main/java/com/zfoo/event/manager/EventBus.java +++ b/event/src/main/java/com/zfoo/event/manager/EventBus.java @@ -59,9 +59,9 @@ public abstract class EventBus { } public static class EventThreadFactory implements ThreadFactory { - private int poolNumber; - private AtomicInteger threadNumber = new AtomicInteger(1); - private ThreadGroup group; + private final int poolNumber; + private final AtomicInteger threadNumber = new AtomicInteger(1); + private final ThreadGroup group; public EventThreadFactory(int poolNumber) { var s = System.getSecurityManager(); @@ -111,7 +111,7 @@ public abstract class EventBus { executors[Math.abs(event.threadId() % EXECUTORS_SIZE)].execute(() -> doSubmit(event, list)); } - public static void asyncExecute(SafeRunnable runnable) { + public static void asyncExecute(Runnable runnable) { execute(RandomUtils.randomInt(EXECUTORS_SIZE), runnable); } @@ -121,8 +121,8 @@ public abstract class EventBus { * @param hashcode * @return */ - public static void execute(int hashcode, SafeRunnable runnable) { - executors[Math.abs(hashcode % EXECUTORS_SIZE)].execute(runnable); + public static void execute(int hashcode, Runnable runnable) { + executors[Math.abs(hashcode % EXECUTORS_SIZE)].execute(SafeRunnable.valueOf(runnable)); } /** diff --git a/event/src/main/java/com/zfoo/event/model/anno/EventReceiver.java b/event/src/main/java/com/zfoo/event/model/anno/EventReceiver.java index 35ba9a6e..b707a5b1 100644 --- a/event/src/main/java/com/zfoo/event/model/anno/EventReceiver.java +++ b/event/src/main/java/com/zfoo/event/model/anno/EventReceiver.java @@ -18,7 +18,7 @@ import java.lang.annotation.*; /** * 接收事件的注解 * - * @author jaysunxiao + * @author godotg * @version 3.0 */ diff --git a/event/src/main/java/com/zfoo/event/model/event/AppStartAfterEvent.java b/event/src/main/java/com/zfoo/event/model/event/AppStartAfterEvent.java index 62f8841e..e340cf15 100644 --- a/event/src/main/java/com/zfoo/event/model/event/AppStartAfterEvent.java +++ b/event/src/main/java/com/zfoo/event/model/event/AppStartAfterEvent.java @@ -21,7 +21,7 @@ import org.springframework.context.event.ApplicationContextEvent; *

* 启动顺序为:AppStartBeforeEvent -> AppStartEvent -> AppStartAfterEvent * - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class AppStartAfterEvent extends ApplicationContextEvent { diff --git a/event/src/main/java/com/zfoo/event/model/event/AppStartBeforeEvent.java b/event/src/main/java/com/zfoo/event/model/event/AppStartBeforeEvent.java index c596963f..c69d7886 100644 --- a/event/src/main/java/com/zfoo/event/model/event/AppStartBeforeEvent.java +++ b/event/src/main/java/com/zfoo/event/model/event/AppStartBeforeEvent.java @@ -21,7 +21,7 @@ import org.springframework.context.event.ApplicationContextEvent; *

* 启动顺序为:AppStartBeforeEvent -> AppStartEvent -> AppStartAfterEvent * - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class AppStartBeforeEvent extends ApplicationContextEvent { diff --git a/event/src/main/java/com/zfoo/event/model/event/AppStartEvent.java b/event/src/main/java/com/zfoo/event/model/event/AppStartEvent.java index 46684048..f7589f99 100644 --- a/event/src/main/java/com/zfoo/event/model/event/AppStartEvent.java +++ b/event/src/main/java/com/zfoo/event/model/event/AppStartEvent.java @@ -21,7 +21,7 @@ import org.springframework.context.event.ApplicationContextEvent; *

* 启动顺序为:AppStartBeforeEvent -> AppStartEvent -> AppStartAfterEvent * - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class AppStartEvent extends ApplicationContextEvent { diff --git a/event/src/main/java/com/zfoo/event/model/event/IEvent.java b/event/src/main/java/com/zfoo/event/model/event/IEvent.java index da11c6a1..e98578aa 100644 --- a/event/src/main/java/com/zfoo/event/model/event/IEvent.java +++ b/event/src/main/java/com/zfoo/event/model/event/IEvent.java @@ -16,7 +16,7 @@ package com.zfoo.event.model.event; import com.zfoo.util.math.RandomUtils; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public interface IEvent { diff --git a/event/src/main/java/com/zfoo/event/model/vo/EnhanceUtils.java b/event/src/main/java/com/zfoo/event/model/vo/EnhanceUtils.java index d339dcc7..f2f1fef0 100644 --- a/event/src/main/java/com/zfoo/event/model/vo/EnhanceUtils.java +++ b/event/src/main/java/com/zfoo/event/model/vo/EnhanceUtils.java @@ -25,7 +25,7 @@ import java.lang.reflect.Method; import java.lang.reflect.Modifier; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public abstract class EnhanceUtils { diff --git a/event/src/main/java/com/zfoo/event/model/vo/EventReceiverDefinition.java b/event/src/main/java/com/zfoo/event/model/vo/EventReceiverDefinition.java index c6ae69b6..94d8ec57 100644 --- a/event/src/main/java/com/zfoo/event/model/vo/EventReceiverDefinition.java +++ b/event/src/main/java/com/zfoo/event/model/vo/EventReceiverDefinition.java @@ -21,7 +21,7 @@ import java.lang.reflect.Method; /** * 动态代理被EventReceiver注解标注的方法,为了避免反射最终会用javassist字节码增强的方法去代理EventReceiverDefinition * - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class EventReceiverDefinition implements IEventReceiver { diff --git a/event/src/main/java/com/zfoo/event/model/vo/IEventReceiver.java b/event/src/main/java/com/zfoo/event/model/vo/IEventReceiver.java index adbe60cd..77d39555 100644 --- a/event/src/main/java/com/zfoo/event/model/vo/IEventReceiver.java +++ b/event/src/main/java/com/zfoo/event/model/vo/IEventReceiver.java @@ -16,7 +16,7 @@ package com.zfoo.event.model.vo; import com.zfoo.event.model.event.IEvent; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public interface IEventReceiver { diff --git a/event/src/main/java/com/zfoo/event/schema/EventDefinitionParser.java b/event/src/main/java/com/zfoo/event/schema/EventDefinitionParser.java index c083fb4c..f514f5c1 100644 --- a/event/src/main/java/com/zfoo/event/schema/EventDefinitionParser.java +++ b/event/src/main/java/com/zfoo/event/schema/EventDefinitionParser.java @@ -22,7 +22,7 @@ import org.springframework.beans.factory.xml.ParserContext; import org.w3c.dom.Element; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class EventDefinitionParser implements BeanDefinitionParser { 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 24163eb8..c1111d2b 100644 --- a/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java +++ b/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java @@ -32,7 +32,7 @@ import java.lang.reflect.Modifier; * 这是一个后置处理器,在boot项目中注册EventContext时,会import导入EventRegisterProcessor这个组件,这是一个后置处理器, * 断点发现 在AbstractAutowireCapableBeanFactory或调用getBeanPostProcessors,这样子每一个后置处理器都会走postProcessAfterInitialization这个方法 * - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class EventRegisterProcessor implements BeanPostProcessor { diff --git a/event/src/main/java/com/zfoo/event/schema/NamespaceHandler.java b/event/src/main/java/com/zfoo/event/schema/NamespaceHandler.java index 97a8683e..bf70d962 100644 --- a/event/src/main/java/com/zfoo/event/schema/NamespaceHandler.java +++ b/event/src/main/java/com/zfoo/event/schema/NamespaceHandler.java @@ -16,7 +16,7 @@ package com.zfoo.event.schema; import org.springframework.beans.factory.xml.NamespaceHandlerSupport; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class NamespaceHandler extends NamespaceHandlerSupport { diff --git a/event/src/test/java/com/zfoo/event/ApplicationTest.java b/event/src/test/java/com/zfoo/event/ApplicationTest.java index 1c3c2e4c..f64a1339 100644 --- a/event/src/test/java/com/zfoo/event/ApplicationTest.java +++ b/event/src/test/java/com/zfoo/event/ApplicationTest.java @@ -20,7 +20,7 @@ import org.junit.Test; import org.springframework.context.support.ClassPathXmlApplicationContext; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ @Ignore diff --git a/event/src/test/java/com/zfoo/event/MyController1.java b/event/src/test/java/com/zfoo/event/MyController1.java index fbbf604f..d7c97ad1 100644 --- a/event/src/test/java/com/zfoo/event/MyController1.java +++ b/event/src/test/java/com/zfoo/event/MyController1.java @@ -19,7 +19,7 @@ import org.slf4j.LoggerFactory; import org.springframework.stereotype.Component; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ @Component diff --git a/event/src/test/java/com/zfoo/event/MyController2.java b/event/src/test/java/com/zfoo/event/MyController2.java index b4992a72..8fc67e11 100644 --- a/event/src/test/java/com/zfoo/event/MyController2.java +++ b/event/src/test/java/com/zfoo/event/MyController2.java @@ -19,7 +19,7 @@ import org.slf4j.LoggerFactory; import org.springframework.stereotype.Component; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ @Component diff --git a/event/src/test/java/com/zfoo/event/MyNoticeEvent.java b/event/src/test/java/com/zfoo/event/MyNoticeEvent.java index a6de3894..4db0bca0 100644 --- a/event/src/test/java/com/zfoo/event/MyNoticeEvent.java +++ b/event/src/test/java/com/zfoo/event/MyNoticeEvent.java @@ -16,7 +16,7 @@ package com.zfoo.event; import com.zfoo.event.model.event.IEvent; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class MyNoticeEvent implements IEvent { 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 9a8341a6..b6c7e635 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 @@ -29,7 +29,6 @@ import com.zfoo.protocol.util.IOUtils; import com.zfoo.protocol.util.JsonUtils; import com.zfoo.protocol.util.StringUtils; import com.zfoo.scheduler.manager.SchedulerBus; -import com.zfoo.util.SafeRunnable; import com.zfoo.util.ThreadUtils; import com.zfoo.util.net.HostAndPort; import io.netty.util.concurrent.FastThreadLocalThread; @@ -361,12 +360,7 @@ public class ZookeeperRegistry implements IRegistry { } catch (Exception e) { // logger.error("zookeeper初始化失败,等待[{}]秒,重新初始化", RETRY_SECONDS, e); - SchedulerBus.schedule(new SafeRunnable() { - @Override - public void doRun() { - initZookeeper(); - } - }, RETRY_SECONDS, TimeUnit.SECONDS); + SchedulerBus.schedule(() -> initZookeeper(), RETRY_SECONDS, TimeUnit.SECONDS); } }); } @@ -522,12 +516,7 @@ public class ZookeeperRegistry implements IRegistry { } if (recheckFlag) { - SchedulerBus.schedule(new SafeRunnable() { - @Override - public void doRun() { - checkConsumer(); - } - }, RETRY_SECONDS, TimeUnit.SECONDS); + SchedulerBus.schedule(() -> checkConsumer(), RETRY_SECONDS, TimeUnit.SECONDS); } } diff --git a/net/src/main/java/com/zfoo/net/task/TaskBus.java b/net/src/main/java/com/zfoo/net/task/TaskBus.java index 4f08a40d..6c01d569 100644 --- a/net/src/main/java/com/zfoo/net/task/TaskBus.java +++ b/net/src/main/java/com/zfoo/net/task/TaskBus.java @@ -73,9 +73,9 @@ public final class TaskBus { } public static class TaskThreadFactory implements ThreadFactory { - private int poolNumber; - private AtomicInteger threadNumber = new AtomicInteger(1); - private ThreadGroup group; + private final int poolNumber; + private final AtomicInteger threadNumber = new AtomicInteger(1); + private final ThreadGroup group; public TaskThreadFactory(int poolNumber) { var s = System.getSecurityManager(); @@ -124,8 +124,8 @@ public final class TaskBus { return Math.abs(executorConsistentHash % EXECUTOR_SIZE); } - public static void execute(int executorConsistentHash, SafeRunnable runnable) { - executors[executorIndex(executorConsistentHash)].execute(runnable); + public static void execute(int executorConsistentHash, Runnable runnable) { + executors[executorIndex(executorConsistentHash)].execute(SafeRunnable.valueOf(runnable)); } // 在task,event,scheduler线程执行的异步请求,请求成功过后依然在相同的线程执行回调任务 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 a4779606..71b9bc96 100644 --- a/net/src/main/java/com/zfoo/net/util/SimpleCache.java +++ b/net/src/main/java/com/zfoo/net/util/SimpleCache.java @@ -20,7 +20,6 @@ import com.zfoo.event.manager.EventBus; import com.zfoo.protocol.collection.CollectionUtils; import com.zfoo.protocol.model.Pair; import com.zfoo.scheduler.manager.SchedulerBus; -import com.zfoo.util.SafeRunnable; import org.checkerframework.checker.nullness.qual.NonNull; import org.checkerframework.checker.nullness.qual.Nullable; @@ -94,31 +93,25 @@ public class SimpleCache { }); - SchedulerBus.scheduleAtFixedRate(new SafeRunnable() { - @Override - public void doRun() { - // 不在任务调度线程中执行耗时任务,因为任务调度线程只有一个线程池 - EventBus.asyncExecute(new SafeRunnable() { - @Override - public void doRun() { - var list = new ArrayList(); - while (!linkedQueue.isEmpty()) { - var key = linkedQueue.poll(); - list.add(key); - if (list.size() >= BATCH_RELOAD_SIZE) { - var result = batchLoadCallback.apply(list); - result.forEach(it -> cache.put(it.getKey(), it.getValue())); - list.clear(); - } - } - - if (CollectionUtils.isNotEmpty(list)) { - var result = batchLoadCallback.apply(list); - result.forEach(it -> cache.put(it.getKey(), it.getValue())); - } + SchedulerBus.scheduleAtFixedRate(() -> { + // 不在任务调度线程中执行耗时任务,因为任务调度线程只有一个线程池 + EventBus.asyncExecute(() -> { + var list = new ArrayList(); + while (!linkedQueue.isEmpty()) { + var key = linkedQueue.poll(); + list.add(key); + if (list.size() >= BATCH_RELOAD_SIZE) { + var result = batchLoadCallback.apply(list); + result.forEach(it -> cache.put(it.getKey(), it.getValue())); + list.clear(); } - }); - } + } + + if (CollectionUtils.isNotEmpty(list)) { + var result = batchLoadCallback.apply(list); + result.forEach(it -> cache.put(it.getKey(), it.getValue())); + } + }); }, refreshDuration, TimeUnit.MILLISECONDS); diff --git a/net/src/main/java/com/zfoo/net/util/SingleCache.java b/net/src/main/java/com/zfoo/net/util/SingleCache.java index bca28574..8ed615fe 100644 --- a/net/src/main/java/com/zfoo/net/util/SingleCache.java +++ b/net/src/main/java/com/zfoo/net/util/SingleCache.java @@ -15,7 +15,6 @@ package com.zfoo.net.util; import com.zfoo.event.manager.EventBus; import com.zfoo.scheduler.util.TimeUtils; -import com.zfoo.util.SafeRunnable; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -62,12 +61,7 @@ public class SingleCache { try { if (now > refreshTime) { refreshTime = now + refreshDuration; - EventBus.asyncExecute(new SafeRunnable() { - @Override - public void doRun() { - cache = supplier.get(); - } - }); + EventBus.asyncExecute(() -> cache = supplier.get()); } } finally { lock.unlock(); diff --git a/net/src/test/java/com/zfoo/net/router/SignalBridgeTest.java b/net/src/test/java/com/zfoo/net/router/SignalBridgeTest.java index ffb23d8e..225291dd 100644 --- a/net/src/test/java/com/zfoo/net/router/SignalBridgeTest.java +++ b/net/src/test/java/com/zfoo/net/router/SignalBridgeTest.java @@ -15,7 +15,6 @@ package com.zfoo.net.router; import com.zfoo.event.manager.EventBus; import com.zfoo.net.router.route.SignalBridge; import com.zfoo.scheduler.util.TimeUtils; -import com.zfoo.util.SafeRunnable; import com.zfoo.util.ThreadUtils; import org.junit.Ignore; import org.junit.Test; @@ -55,9 +54,9 @@ public class SignalBridgeTest { var countDownLatch = new CountDownLatch(executorSize); for (var i = 0; i < executorSize; i++) { - EventBus.execute(i, new SafeRunnable() { + EventBus.execute(i, new Runnable() { @Override - public void doRun() { + public void run() { addAndRemoveArray(); countDownLatch.countDown(); } @@ -73,9 +72,9 @@ public class SignalBridgeTest { var countDownLatch = new CountDownLatch(executorSize); for (int i = 0; i < executorSize; i++) { - EventBus.execute(i, new SafeRunnable() { + EventBus.execute(i, new Runnable() { @Override - public void doRun() { + public void run() { addAndRemoveMap(); countDownLatch.countDown(); } 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 d6d64c5a..af5573d0 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 @@ -20,7 +20,6 @@ import com.zfoo.orm.model.vo.EntityDef; import com.zfoo.protocol.exception.ExceptionUtils; import com.zfoo.scheduler.manager.SchedulerBus; import com.zfoo.scheduler.util.TimeUtils; -import com.zfoo.util.SafeRunnable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.scheduling.support.CronExpression; @@ -43,7 +42,7 @@ public class CronOrmPersister extends AbstractOrmPersister { /** * cron表达式 */ - private CronExpression cronExpression; + private final CronExpression cronExpression; public CronOrmPersister(EntityDef entityDef, EntityCaches entityCaches) { @@ -75,18 +74,12 @@ public class CronOrmPersister extends AbstractOrmPersister { } if (!OrmContext.isStop()) { - SchedulerBus.schedule(new SafeRunnable() { - @Override - public void doRun() { - if (!OrmContext.isStop()) { - EventBus.execute(entityDef.getClazz().hashCode(), new SafeRunnable() { - @Override - public void doRun() { - entityCaches.persistAll(); - schedulePersist(); - } - }); - } + SchedulerBus.schedule(() -> { + if (!OrmContext.isStop()) { + EventBus.execute(entityDef.getClazz().hashCode(), () -> { + entityCaches.persistAll(); + schedulePersist(); + }); } }, delay, TimeUnit.MILLISECONDS); } 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 393fd816..cafc9121 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 @@ -19,7 +19,6 @@ import com.zfoo.orm.model.cache.EntityCaches; import com.zfoo.orm.model.vo.EntityDef; import com.zfoo.protocol.util.StringUtils; import com.zfoo.scheduler.manager.SchedulerBus; -import com.zfoo.util.SafeRunnable; import java.util.concurrent.TimeUnit; @@ -32,7 +31,7 @@ public class TimeOrmPersister extends AbstractOrmPersister { /** * 执行的频率 */ - private long rate; + private final long rate; public TimeOrmPersister(EntityDef entityDef, EntityCaches entityCaches) { super(entityDef, entityCaches); @@ -44,17 +43,9 @@ public class TimeOrmPersister extends AbstractOrmPersister { @Override public void start() { - SchedulerBus.scheduleAtFixedRate(new SafeRunnable() { - @Override - public void doRun() { - if (!OrmContext.isStop()) { - EventBus.execute(entityDef.getClazz().hashCode(), new SafeRunnable() { - @Override - public void doRun() { - entityCaches.persistAll(); - } - }); - } + SchedulerBus.scheduleAtFixedRate(() -> { + if (!OrmContext.isStop()) { + EventBus.execute(entityDef.getClazz().hashCode(), () -> entityCaches.persistAll()); } }, rate, TimeUnit.MILLISECONDS); } diff --git a/orm/src/test/java/com/zfoo/orm/lpmap/ConcurrentFileChannelMapTest.java b/orm/src/test/java/com/zfoo/orm/lpmap/ConcurrentFileChannelMapTest.java index a167caed..ee9c6f9e 100644 --- a/orm/src/test/java/com/zfoo/orm/lpmap/ConcurrentFileChannelMapTest.java +++ b/orm/src/test/java/com/zfoo/orm/lpmap/ConcurrentFileChannelMapTest.java @@ -15,7 +15,6 @@ package com.zfoo.orm.lpmap; import com.zfoo.event.manager.EventBus; import com.zfoo.orm.lpmap.model.MyPacket; import com.zfoo.protocol.ProtocolManager; -import com.zfoo.util.SafeRunnable; import org.junit.Assert; import org.junit.Ignore; import org.junit.Test; @@ -42,9 +41,9 @@ public class ConcurrentFileChannelMapTest { var countdown = new CountDownLatch(EventBus.EXECUTORS_SIZE); for (int i = 0; i < EventBus.EXECUTORS_SIZE; i++) { - EventBus.asyncExecute(new SafeRunnable() { + EventBus.asyncExecute(new Runnable() { @Override - public void doRun() { + public void run() { var key = atomicInt.getAndIncrement(); while (key < count) { var myPacket = MyPacket.valueOf(key, String.valueOf(key)); diff --git a/orm/src/test/java/com/zfoo/orm/lpmap/ConcurrentHeapMapTest.java b/orm/src/test/java/com/zfoo/orm/lpmap/ConcurrentHeapMapTest.java index ce9895d5..69c45ed2 100644 --- a/orm/src/test/java/com/zfoo/orm/lpmap/ConcurrentHeapMapTest.java +++ b/orm/src/test/java/com/zfoo/orm/lpmap/ConcurrentHeapMapTest.java @@ -15,7 +15,6 @@ package com.zfoo.orm.lpmap; import com.zfoo.event.manager.EventBus; import com.zfoo.orm.lpmap.model.MyPacket; import com.zfoo.protocol.ProtocolManager; -import com.zfoo.util.SafeRunnable; import org.junit.Assert; import org.junit.Ignore; import org.junit.Test; @@ -59,17 +58,14 @@ public class ConcurrentHeapMapTest { var countdown = new CountDownLatch(EventBus.EXECUTORS_SIZE); for (int i = 0; i < EventBus.EXECUTORS_SIZE; i++) { - EventBus.asyncExecute(new SafeRunnable() { - @Override - public void doRun() { - var key = atomicInt.getAndIncrement(); - while (key < count) { - var myPacket = MyPacket.valueOf(key, String.valueOf(key)); - map.put(key, myPacket); - key = atomicInt.getAndIncrement(); - } - countdown.countDown(); + EventBus.asyncExecute(() -> { + var key = atomicInt.getAndIncrement(); + while (key < count) { + var myPacket = MyPacket.valueOf(key, String.valueOf(key)); + map.put(key, myPacket); + key = atomicInt.getAndIncrement(); } + countdown.countDown(); }); } countdown.await(); diff --git a/scheduler/src/main/java/com/zfoo/scheduler/SchedulerContext.java b/scheduler/src/main/java/com/zfoo/scheduler/SchedulerContext.java index 9d32b5a7..f1a08f2a 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/SchedulerContext.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/SchedulerContext.java @@ -34,7 +34,7 @@ import java.lang.reflect.Modifier; import java.util.concurrent.ScheduledExecutorService; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class SchedulerContext implements ApplicationListener, Ordered { diff --git a/scheduler/src/main/java/com/zfoo/scheduler/manager/SchedulerBus.java b/scheduler/src/main/java/com/zfoo/scheduler/manager/SchedulerBus.java index fb4e3115..9d7060a7 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/manager/SchedulerBus.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/manager/SchedulerBus.java @@ -72,7 +72,7 @@ public abstract class SchedulerBus { public static class SchedulerThreadFactory implements ThreadFactory { - private int poolNumber; + private final int poolNumber; private final AtomicInteger threadNumber = new AtomicInteger(1); private final ThreadGroup group; @@ -166,24 +166,24 @@ public abstract class SchedulerBus { /** * 不断执行的周期循环任务 */ - public static void scheduleAtFixedRate(SafeRunnable runnable, long period, TimeUnit unit) { + public static void scheduleAtFixedRate(Runnable runnable, long period, TimeUnit unit) { if (SchedulerContext.isStop()) { return; } - executor.scheduleAtFixedRate(runnable, 0, period, unit); + executor.scheduleAtFixedRate(SafeRunnable.valueOf(runnable), 0, period, unit); } /** * 固定延迟执行的任务 */ - public static void schedule(SafeRunnable runnable, long delay, TimeUnit unit) { + public static void schedule(Runnable runnable, long delay, TimeUnit unit) { if (SchedulerContext.isStop()) { return; } - executor.schedule(runnable, delay, unit); + executor.schedule(SafeRunnable.valueOf(runnable), delay, unit); } /** diff --git a/scheduler/src/main/java/com/zfoo/scheduler/model/StopWatch.java b/scheduler/src/main/java/com/zfoo/scheduler/model/StopWatch.java index d9af4534..487a2b79 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/model/StopWatch.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/model/StopWatch.java @@ -18,12 +18,12 @@ import java.math.BigDecimal; import java.math.RoundingMode; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class StopWatch { - private long startTime = TimeUtils.currentTimeMillis(); + private final long startTime = TimeUtils.currentTimeMillis(); public long cost() { return TimeUtils.currentTimeMillis() - startTime; diff --git a/scheduler/src/main/java/com/zfoo/scheduler/model/anno/Scheduler.java b/scheduler/src/main/java/com/zfoo/scheduler/model/anno/Scheduler.java index 4e61defd..430a451e 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/model/anno/Scheduler.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/model/anno/Scheduler.java @@ -16,7 +16,7 @@ package com.zfoo.scheduler.model.anno; import java.lang.annotation.*; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ @Documented diff --git a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/EnhanceUtils.java b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/EnhanceUtils.java index 5772917e..6d8f202f 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/EnhanceUtils.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/EnhanceUtils.java @@ -24,7 +24,7 @@ import java.lang.reflect.Method; import java.lang.reflect.Modifier; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public abstract class EnhanceUtils { diff --git a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/IScheduler.java b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/IScheduler.java index 6361288b..8259792c 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/IScheduler.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/IScheduler.java @@ -14,7 +14,7 @@ package com.zfoo.scheduler.model.vo; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public interface IScheduler { diff --git a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/ReflectScheduler.java b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/ReflectScheduler.java index d29ceb2d..c45566fd 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/ReflectScheduler.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/ReflectScheduler.java @@ -20,7 +20,7 @@ import java.lang.reflect.Method; /** * 动态代理被Scheduler注解标注的方法,为了避免反射最终会用javassist字节码增强的方法去代理ReflectScheduler * - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class ReflectScheduler implements IScheduler { diff --git a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/RunnableScheduler.java b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/RunnableScheduler.java index 2eaf76da..9cb233ae 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/RunnableScheduler.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/RunnableScheduler.java @@ -14,7 +14,7 @@ package com.zfoo.scheduler.model.vo; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class RunnableScheduler implements IScheduler { diff --git a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/SchedulerDefinition.java b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/SchedulerDefinition.java index ca07d211..e94d1be5 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/model/vo/SchedulerDefinition.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/model/vo/SchedulerDefinition.java @@ -23,7 +23,7 @@ import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Method; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class SchedulerDefinition { diff --git a/scheduler/src/main/java/com/zfoo/scheduler/schema/NamespaceHandler.java b/scheduler/src/main/java/com/zfoo/scheduler/schema/NamespaceHandler.java index 8a125e41..62f856df 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/schema/NamespaceHandler.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/schema/NamespaceHandler.java @@ -16,7 +16,7 @@ package com.zfoo.scheduler.schema; import org.springframework.beans.factory.xml.NamespaceHandlerSupport; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class NamespaceHandler extends NamespaceHandlerSupport { 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 92d76813..14c373eb 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerDefinitionParser.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/schema/SchedulerDefinitionParser.java @@ -31,7 +31,7 @@ import org.w3c.dom.Element; * 从而可以得出结论:在基于zfoo的SpringBoot工程中,由于压根不解析xml对象,注册bean也是通过@Configuration配置类注册bean的, * 因此这些实现BeanDefinitionParser接口的parse方法是都不会执行的。 * - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class SchedulerDefinitionParser implements BeanDefinitionParser { 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 76ca47b4..76bca3bf 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/util/TimeUtils.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/util/TimeUtils.java @@ -26,7 +26,7 @@ import java.util.Date; import java.util.TimeZone; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public abstract class TimeUtils { diff --git a/scheduler/src/test/java/com/zfoo/scheduler/ApplicationTest.java b/scheduler/src/test/java/com/zfoo/scheduler/ApplicationTest.java index 1e968201..2a7fdc51 100644 --- a/scheduler/src/test/java/com/zfoo/scheduler/ApplicationTest.java +++ b/scheduler/src/test/java/com/zfoo/scheduler/ApplicationTest.java @@ -22,7 +22,7 @@ import org.springframework.context.support.ClassPathXmlApplicationContext; * cron(译为克龙)代表100万年,是英文单词中最大的时间单位。 * google(译为古戈尔)代表10的100次方,足够穷尽宇宙万物 * - * @author jaysunxiao + * @author godotg * @version 3.0 */ diff --git a/scheduler/src/test/java/com/zfoo/scheduler/SchedulerController.java b/scheduler/src/test/java/com/zfoo/scheduler/SchedulerController.java index adc9cc49..9c5e4482 100644 --- a/scheduler/src/test/java/com/zfoo/scheduler/SchedulerController.java +++ b/scheduler/src/test/java/com/zfoo/scheduler/SchedulerController.java @@ -19,7 +19,7 @@ import org.slf4j.LoggerFactory; import org.springframework.stereotype.Component; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ @Component diff --git a/scheduler/src/test/java/com/zfoo/scheduler/util/TimeUtilsTest.java b/scheduler/src/test/java/com/zfoo/scheduler/util/TimeUtilsTest.java index 302a55a4..e7c02c63 100644 --- a/scheduler/src/test/java/com/zfoo/scheduler/util/TimeUtilsTest.java +++ b/scheduler/src/test/java/com/zfoo/scheduler/util/TimeUtilsTest.java @@ -25,7 +25,7 @@ import java.time.temporal.TemporalAdjusters; import java.util.Date; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class TimeUtilsTest { diff --git a/util/src/main/java/com/zfoo/util/SafeRunnable.java b/util/src/main/java/com/zfoo/util/SafeRunnable.java index 61d8f654..8dd28929 100644 --- a/util/src/main/java/com/zfoo/util/SafeRunnable.java +++ b/util/src/main/java/com/zfoo/util/SafeRunnable.java @@ -7,14 +7,25 @@ import org.slf4j.LoggerFactory; * @author godotg * @version 3.0 */ -public abstract class SafeRunnable implements Runnable { +public class SafeRunnable implements Runnable { private static final Logger logger = LoggerFactory.getLogger(SafeRunnable.class); + private Runnable runnable; + + private SafeRunnable() { + } + + public static SafeRunnable valueOf(Runnable runnable) { + var run = new SafeRunnable(); + run.runnable = runnable; + return run; + } + @Override public void run() { try { - doRun(); + runnable.run(); } catch (Exception e) { logger.error("未知exception异常", e); } catch (Throwable t) { @@ -22,6 +33,4 @@ public abstract class SafeRunnable implements Runnable { } } - public abstract void doRun(); - }