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();
-
}