From 070000fbe4164520fcdc3f8b9d7feca543548b19 Mon Sep 17 00:00:00 2001 From: godotg Date: Wed, 10 Apr 2024 18:05:00 +0800 Subject: [PATCH] feat[event]: do not init event executors when not use AsyncThread --- .../java/com/zfoo/event/manager/EventBus.java | 60 +-------------- .../zfoo/event/manager/EventExecutors.java | 77 +++++++++++++++++++ .../main/java/com/zfoo/net/task/TaskBus.java | 4 +- .../com/zfoo/net/router/SignalBridgeTest.java | 2 +- .../java/com/zfoo/orm/cache/EntityCache.java | 2 +- .../orm/cache/persister/CronOrmPersister.java | 2 +- .../orm/cache/persister/TimeOrmPersister.java | 2 +- 7 files changed, 87 insertions(+), 62 deletions(-) create mode 100644 event/src/main/java/com/zfoo/event/manager/EventExecutors.java 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 51668668..ca822626 100644 --- a/event/src/main/java/com/zfoo/event/manager/EventBus.java +++ b/event/src/main/java/com/zfoo/event/manager/EventBus.java @@ -15,12 +15,8 @@ package com.zfoo.event.manager; import com.zfoo.event.enhance.IEventReceiver; import com.zfoo.event.model.IEvent; import com.zfoo.protocol.collection.CollectionUtils; -import com.zfoo.protocol.collection.concurrent.CopyOnWriteHashMapLongObject; -import com.zfoo.protocol.util.AssertionUtils; import com.zfoo.protocol.util.RandomUtils; -import com.zfoo.protocol.util.StringUtils; import com.zfoo.protocol.util.ThreadUtils; -import io.netty.util.concurrent.FastThreadLocalThread; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -28,11 +24,6 @@ import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.concurrent.Executor; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.ThreadFactory; -import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiConsumer; import java.util.function.Consumer; @@ -43,15 +34,6 @@ public abstract class EventBus { private static final Logger logger = LoggerFactory.getLogger(EventBus.class); - /** - * EN: The size of the thread pool. Event's thread pool is often used to do time-consuming operations, so set it a little bigger - * CN: 线程池的大小. event的线程池经常用来做一些耗时的操作,所以要设置大一点 - */ - private static final int EXECUTORS_SIZE = Math.max(Runtime.getRuntime().availableProcessors(), 4) * 2 + 1; - - private static final ExecutorService[] executors = new ExecutorService[EXECUTORS_SIZE]; - - private static final CopyOnWriteHashMapLongObject threadMap = new CopyOnWriteHashMapLongObject<>(EXECUTORS_SIZE); /** * event mapping */ @@ -68,37 +50,6 @@ public abstract class EventBus { public static BiConsumer exceptionHandler = null; public static Consumer noEventReceiverHandler = null; - static { - for (int i = 0; i < executors.length; i++) { - var namedThreadFactory = new EventThreadFactory(i); - var executor = Executors.newSingleThreadExecutor(namedThreadFactory); - executors[i] = executor; - } - } - - public static class EventThreadFactory implements ThreadFactory { - private final int poolNumber; - private final AtomicInteger threadNumber = new AtomicInteger(1); - private final ThreadGroup group; - - public EventThreadFactory(int poolNumber) { - this.group = Thread.currentThread().getThreadGroup(); - this.poolNumber = poolNumber; - } - - @Override - public Thread newThread(Runnable runnable) { - var threadName = StringUtils.format("event-p{}-t{}", poolNumber + 1, threadNumber.getAndIncrement()); - var thread = new FastThreadLocalThread(group, runnable, threadName); - thread.setDaemon(false); - thread.setPriority(Thread.NORM_PRIORITY); - thread.setUncaughtExceptionHandler((t, e) -> logger.error(t.toString(), e)); - var executor = executors[poolNumber]; - AssertionUtils.notNull(executor); - threadMap.put(thread.getId(), executor); - return thread; - } - } /** * Publish the event @@ -120,7 +71,7 @@ public abstract class EventBus { for (var receiver : receivers) { switch (receiver.bus()) { case CurrentThread -> doReceiver(receiver, event); - case AsyncThread -> execute(event.executorHash(), () -> doReceiver(receiver, event)); + case AsyncThread -> asyncExecute(event.executorHash(), () -> doReceiver(receiver, event)); // case VirtualThread -> Thread.ofVirtual().name("virtual-on" + clazz.getSimpleName()).start(() -> doReceiver(receiver, event)); case ManualThread -> manualThreadHandler.accept(receiver, event); } @@ -139,14 +90,14 @@ public abstract class EventBus { } public static void asyncExecute(Runnable runnable) { - execute(RandomUtils.randomInt(), runnable); + asyncExecute(RandomUtils.randomInt(), runnable); } /** * Use the event thread specified by the hashcode to execute the task */ - public static void execute(int executorHash, Runnable runnable) { - executors[Math.abs(executorHash % EXECUTORS_SIZE)].execute(ThreadUtils.safeRunnable(runnable)); + public static void asyncExecute(int executorHash, Runnable runnable) { + EventExecutors.execute(executorHash, ThreadUtils.safeRunnable(runnable)); } /** @@ -156,9 +107,6 @@ public abstract class EventBus { receiverMap.computeIfAbsent(eventType, it -> new ArrayList<>(1)).add(receiver); } - public static Executor threadExecutor(long currentThreadId) { - return threadMap.getPrimitive(currentThreadId); - } } diff --git a/event/src/main/java/com/zfoo/event/manager/EventExecutors.java b/event/src/main/java/com/zfoo/event/manager/EventExecutors.java new file mode 100644 index 00000000..86930aeb --- /dev/null +++ b/event/src/main/java/com/zfoo/event/manager/EventExecutors.java @@ -0,0 +1,77 @@ +package com.zfoo.event.manager; + +import com.zfoo.protocol.collection.concurrent.CopyOnWriteHashMapLongObject; +import com.zfoo.protocol.util.AssertionUtils; +import com.zfoo.protocol.util.StringUtils; +import com.zfoo.protocol.util.ThreadUtils; +import io.netty.util.concurrent.FastThreadLocalThread; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.concurrent.Executor; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * @author godotg + */ +public abstract class EventExecutors { + + private static final Logger logger = LoggerFactory.getLogger(EventExecutors.class); + + /** + * EN: The size of the thread pool. Event's thread pool is often used to do time-consuming operations, so set it a little bigger + * CN: 线程池的大小. event的线程池经常用来做一些耗时的操作,所以要设置大一点 + */ + private static final int EXECUTORS_SIZE = Math.max(Runtime.getRuntime().availableProcessors(), 4) * 2 + 1; + + private static final ExecutorService[] executors = new ExecutorService[EXECUTORS_SIZE]; + + private static final CopyOnWriteHashMapLongObject threadMap = new CopyOnWriteHashMapLongObject<>(EXECUTORS_SIZE); + + + static { + for (int i = 0; i < executors.length; i++) { + var namedThreadFactory = new EventThreadFactory(i); + var executor = Executors.newSingleThreadExecutor(namedThreadFactory); + executors[i] = executor; + } + } + + public static class EventThreadFactory implements ThreadFactory { + private final int poolNumber; + private final AtomicInteger threadNumber = new AtomicInteger(1); + private final ThreadGroup group; + + public EventThreadFactory(int poolNumber) { + this.group = Thread.currentThread().getThreadGroup(); + this.poolNumber = poolNumber; + } + + @Override + public Thread newThread(Runnable runnable) { + var threadName = StringUtils.format("event-p{}-t{}", poolNumber + 1, threadNumber.getAndIncrement()); + var thread = new FastThreadLocalThread(group, runnable, threadName); + thread.setDaemon(false); + thread.setPriority(Thread.NORM_PRIORITY); + thread.setUncaughtExceptionHandler((t, e) -> logger.error(t.toString(), e)); + var executor = executors[poolNumber]; + AssertionUtils.notNull(executor); + threadMap.put(thread.getId(), executor); + return thread; + } + } + + /** + * Use the event thread specified by the hashcode to execute the task + */ + public static void execute(int executorHash, Runnable runnable) { + executors[Math.abs(executorHash % EXECUTORS_SIZE)].execute(ThreadUtils.safeRunnable(runnable)); + } + + public static Executor threadExecutor(long currentThreadId) { + return threadMap.getPrimitive(currentThreadId); + } +} 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 417b2d37..fd34c4cc 100644 --- a/net/src/main/java/com/zfoo/net/task/TaskBus.java +++ b/net/src/main/java/com/zfoo/net/task/TaskBus.java @@ -13,7 +13,7 @@ package com.zfoo.net.task; -import com.zfoo.event.manager.EventBus; +import com.zfoo.event.manager.EventExecutors; import com.zfoo.net.NetContext; import com.zfoo.protocol.collection.concurrent.CopyOnWriteHashMapLongObject; import com.zfoo.protocol.util.AssertionUtils; @@ -126,7 +126,7 @@ public final class TaskBus { return taskExecutor; } - var eventExecutor = EventBus.threadExecutor(threadId); + var eventExecutor = EventExecutors.threadExecutor(threadId); if (eventExecutor != null) { return eventExecutor; } 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 dedb4024..69f6b574 100644 --- a/net/src/test/java/com/zfoo/net/router/SignalBridgeTest.java +++ b/net/src/test/java/com/zfoo/net/router/SignalBridgeTest.java @@ -45,7 +45,7 @@ public class SignalBridgeTest { var countDownLatch = new CountDownLatch(executorSize); for (var i = 0; i < executorSize; i++) { - EventBus.execute(i, new Runnable() { + EventBus.asyncExecute(i, new Runnable() { @Override public void run() { addAndRemoveArray(); diff --git a/orm/src/main/java/com/zfoo/orm/cache/EntityCache.java b/orm/src/main/java/com/zfoo/orm/cache/EntityCache.java index 14abbe29..9416af25 100644 --- a/orm/src/main/java/com/zfoo/orm/cache/EntityCache.java +++ b/orm/src/main/java/com/zfoo/orm/cache/EntityCache.java @@ -69,7 +69,7 @@ public class EntityCache, E extends IEntity> imple var entity = pnode.getEntity(); @SuppressWarnings("unchecked") var entityClass = (Class) entityDef.getClazz(); - EventBus.execute(entityClass.hashCode(), new Runnable() { + EventBus.asyncExecute(entityClass.hashCode(), new Runnable() { @Override public void run() { var collection = OrmContext.getOrmManager().getCollection(entityClass); diff --git a/orm/src/main/java/com/zfoo/orm/cache/persister/CronOrmPersister.java b/orm/src/main/java/com/zfoo/orm/cache/persister/CronOrmPersister.java index 4cb6757f..47938ee4 100644 --- a/orm/src/main/java/com/zfoo/orm/cache/persister/CronOrmPersister.java +++ b/orm/src/main/java/com/zfoo/orm/cache/persister/CronOrmPersister.java @@ -75,7 +75,7 @@ public class CronOrmPersister extends AbstractOrmPersister { if (!OrmContext.isStop()) { SchedulerBus.schedule(() -> { if (!OrmContext.isStop()) { - EventBus.execute(entityDef.getClazz().hashCode(), () -> { + EventBus.asyncExecute(entityDef.getClazz().hashCode(), () -> { entityCaches.persistAll(); schedulePersist(); }); diff --git a/orm/src/main/java/com/zfoo/orm/cache/persister/TimeOrmPersister.java b/orm/src/main/java/com/zfoo/orm/cache/persister/TimeOrmPersister.java index 723ad5d1..4efe6cae 100644 --- a/orm/src/main/java/com/zfoo/orm/cache/persister/TimeOrmPersister.java +++ b/orm/src/main/java/com/zfoo/orm/cache/persister/TimeOrmPersister.java @@ -44,7 +44,7 @@ public class TimeOrmPersister extends AbstractOrmPersister { public void start() { SchedulerBus.scheduleAtFixedRate(() -> { if (!OrmContext.isStop()) { - EventBus.execute(entityDef.getClazz().hashCode(), () -> entityCaches.persistAll()); + EventBus.asyncExecute(entityDef.getClazz().hashCode(), () -> entityCaches.persistAll()); } }, rate, TimeUnit.MILLISECONDS); }