From 949cee885c6f712f0a2b91d5226537f507063342 Mon Sep 17 00:00:00 2001 From: godotg Date: Sun, 11 Dec 2022 18:57:26 +0800 Subject: [PATCH] del[net]: remove difficult and infrequently used task-dispatcher configurations --- .../zfoo/net/config/model/ProviderConfig.java | 15 +---- .../main/java/com/zfoo/net/router/Router.java | 4 +- .../com/zfoo/net/schema/NamespaceHandler.java | 2 +- .../zfoo/net/schema/NetDefinitionParser.java | 3 +- .../task/{model => }/PacketReceiverTask.java | 5 +- .../main/java/com/zfoo/net/task/TaskBus.java | 36 ++++++------ .../task/dispatcher/AbstractTaskDispatch.java | 37 ------------- .../ConsistentHashTaskDispatch.java | 55 ------------------- .../net/task/dispatcher/ITaskDispatch.java | 29 ---------- .../task/dispatcher/RandomTaskDispatch.java | 40 -------------- .../dispatcher/SessionIdTaskDispatch.java | 43 --------------- net/src/main/resources/net-1.0.xsd | 1 - .../resources/provider/provider_config.xml | 2 +- 13 files changed, 28 insertions(+), 244 deletions(-) rename net/src/main/java/com/zfoo/net/task/{model => }/PacketReceiverTask.java (96%) delete mode 100644 net/src/main/java/com/zfoo/net/task/dispatcher/AbstractTaskDispatch.java delete mode 100644 net/src/main/java/com/zfoo/net/task/dispatcher/ConsistentHashTaskDispatch.java delete mode 100644 net/src/main/java/com/zfoo/net/task/dispatcher/ITaskDispatch.java delete mode 100644 net/src/main/java/com/zfoo/net/task/dispatcher/RandomTaskDispatch.java delete mode 100644 net/src/main/java/com/zfoo/net/task/dispatcher/SessionIdTaskDispatch.java diff --git a/net/src/main/java/com/zfoo/net/config/model/ProviderConfig.java b/net/src/main/java/com/zfoo/net/config/model/ProviderConfig.java index 757a086d..d5ef24a9 100644 --- a/net/src/main/java/com/zfoo/net/config/model/ProviderConfig.java +++ b/net/src/main/java/com/zfoo/net/config/model/ProviderConfig.java @@ -26,12 +26,7 @@ import java.util.Objects; */ public class ProviderConfig { - public static transient final int DEFAULT_PORT = 12400; - - /** - * 对应于ITaskDispatch - */ - private String taskDispatch; + public static final int DEFAULT_PORT = 12400; private String thread; @@ -55,14 +50,6 @@ public class ProviderConfig { return HostAndPort.valueOf(address); } - public String getTaskDispatch() { - return taskDispatch; - } - - public void setTaskDispatch(String taskDispatch) { - this.taskDispatch = taskDispatch; - } - public String getThread() { return thread; } diff --git a/net/src/main/java/com/zfoo/net/router/Router.java b/net/src/main/java/com/zfoo/net/router/Router.java index 62d76285..cc65c7d6 100644 --- a/net/src/main/java/com/zfoo/net/router/Router.java +++ b/net/src/main/java/com/zfoo/net/router/Router.java @@ -34,8 +34,8 @@ import com.zfoo.net.router.route.PacketBus; import com.zfoo.net.router.route.SignalBridge; import com.zfoo.net.session.model.AttributeType; import com.zfoo.net.session.model.Session; +import com.zfoo.net.task.PacketReceiverTask; import com.zfoo.net.task.TaskBus; -import com.zfoo.net.task.model.PacketReceiverTask; import com.zfoo.protocol.IPacket; import com.zfoo.protocol.exception.ExceptionUtils; import com.zfoo.protocol.util.JsonUtils; @@ -151,7 +151,7 @@ public class Router implements IRouter { // 正常发送消息的接收,把客户端的业务请求包装下到路由策略指定的线程进行业务处理 // 注意:像客户端以asyncAsk发送请求,在服务器处理完后返回结果,在请求方也是进入这个receive方法,但是attachment不为空,会提前return掉不会走到这 - TaskBus.submit(new PacketReceiverTask(session, packet, attachment)); + TaskBus.dispatch(new PacketReceiverTask(session, packet, attachment)); } @Override diff --git a/net/src/main/java/com/zfoo/net/schema/NamespaceHandler.java b/net/src/main/java/com/zfoo/net/schema/NamespaceHandler.java index 43ed650d..b6c0195f 100644 --- a/net/src/main/java/com/zfoo/net/schema/NamespaceHandler.java +++ b/net/src/main/java/com/zfoo/net/schema/NamespaceHandler.java @@ -16,7 +16,7 @@ package com.zfoo.net.schema; import org.springframework.beans.factory.xml.NamespaceHandlerSupport; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class NamespaceHandler extends NamespaceHandlerSupport { diff --git a/net/src/main/java/com/zfoo/net/schema/NetDefinitionParser.java b/net/src/main/java/com/zfoo/net/schema/NetDefinitionParser.java index 094b26b6..55a64fc7 100644 --- a/net/src/main/java/com/zfoo/net/schema/NetDefinitionParser.java +++ b/net/src/main/java/com/zfoo/net/schema/NetDefinitionParser.java @@ -32,7 +32,7 @@ import org.springframework.beans.factory.xml.ParserContext; import org.w3c.dom.Element; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public class NetDefinitionParser implements BeanDefinitionParser { @@ -179,7 +179,6 @@ public class NetDefinitionParser implements BeanDefinitionParser { var clazz = ProviderConfig.class; var builder = BeanDefinitionBuilder.rootBeanDefinition(clazz); - resolvePlaceholder("task-dispatch", "taskDispatch", builder, element, parserContext); resolvePlaceholder("thread", "thread", builder, element, parserContext); resolvePlaceholder("address", "address", builder, element, parserContext); diff --git a/net/src/main/java/com/zfoo/net/task/model/PacketReceiverTask.java b/net/src/main/java/com/zfoo/net/task/PacketReceiverTask.java similarity index 96% rename from net/src/main/java/com/zfoo/net/task/model/PacketReceiverTask.java rename to net/src/main/java/com/zfoo/net/task/PacketReceiverTask.java index 41faefcd..94b393fe 100644 --- a/net/src/main/java/com/zfoo/net/task/model/PacketReceiverTask.java +++ b/net/src/main/java/com/zfoo/net/task/PacketReceiverTask.java @@ -1,6 +1,5 @@ /* * 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 * @@ -11,7 +10,7 @@ * See the License for the specific language governing permissions and limitations under the License. */ -package com.zfoo.net.task.model; +package com.zfoo.net.task; import com.zfoo.net.NetContext; import com.zfoo.net.router.attachment.IAttachment; @@ -19,7 +18,7 @@ import com.zfoo.net.session.model.Session; import com.zfoo.protocol.IPacket; /** - * @author jaysunxiao + * @author godotg * @version 3.0 */ public final class PacketReceiverTask implements Runnable { 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 960d9ce8..ba6a927a 100644 --- a/net/src/main/java/com/zfoo/net/task/TaskBus.java +++ b/net/src/main/java/com/zfoo/net/task/TaskBus.java @@ -15,9 +15,7 @@ package com.zfoo.net.task; import com.zfoo.event.manager.EventBus; import com.zfoo.net.NetContext; -import com.zfoo.net.task.dispatcher.AbstractTaskDispatch; -import com.zfoo.net.task.dispatcher.ITaskDispatch; -import com.zfoo.net.task.model.PacketReceiverTask; +import com.zfoo.net.session.model.AttributeType; import com.zfoo.protocol.collection.concurrent.CopyOnWriteHashMapLongObject; import com.zfoo.protocol.util.AssertionUtils; import com.zfoo.protocol.util.StringUtils; @@ -29,8 +27,10 @@ import io.netty.util.concurrent.FastThreadLocalThread; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.Map; -import java.util.concurrent.*; +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; /** @@ -46,9 +46,6 @@ public final class TaskBus { // 线程池的大小,也可以通过provider thread配置指定 public static final int EXECUTOR_SIZE; - private static final ITaskDispatch taskDispatch; - - /** * 使用不同的线程池,让线程池之间实现隔离,互不影响 */ @@ -58,11 +55,7 @@ public final class TaskBus { var localConfig = NetContext.getConfigManager().getLocalConfig(); var providerConfig = localConfig.getProvider(); - taskDispatch = AbstractTaskDispatch.valueOf(providerConfig == null ? "consistent-hash" : providerConfig.getTaskDispatch()); - - EXECUTOR_SIZE = (providerConfig == null || StringUtils.isBlank(providerConfig.getThread())) - ? (Runtime.getRuntime().availableProcessors() + 1) - : Integer.parseInt(providerConfig.getThread()); + EXECUTOR_SIZE = (providerConfig == null || StringUtils.isBlank(providerConfig.getThread())) ? (Runtime.getRuntime().availableProcessors() + 1) : Integer.parseInt(providerConfig.getThread()); executors = new ExecutorService[EXECUTOR_SIZE]; for (int i = 0; i < executors.length; i++) { @@ -116,9 +109,20 @@ public final class TaskBus { * GatewayAttachment:默认是executorConsistentHash等于用户活玩家的uid,也可以通过IGatewayLoadBalancer接口指定 * SignalAttachment:executorConsistentHash通过IRouter和IConsumer的argument参数指定 */ - public static void submit(PacketReceiverTask task) { - // 里面会看到是:其中一致性hash是根据附加包记录的hashId进行选择哪个线程进行业务处理 - taskDispatch.getExecutor(executors, task).execute(task); + public static void dispatch(PacketReceiverTask task) { + var attachment = task.getAttachment(); + + if (attachment == null) { + var session = task.getSession(); + var uid = session.getAttribute(AttributeType.UID); + if (uid == null) { + execute((int) session.getSid(), task); + } else { + execute((int) uid, task); + } + } else { + execute(attachment.executorConsistentHash(), task); + } } public static int executorIndex(int executorConsistentHash) { diff --git a/net/src/main/java/com/zfoo/net/task/dispatcher/AbstractTaskDispatch.java b/net/src/main/java/com/zfoo/net/task/dispatcher/AbstractTaskDispatch.java deleted file mode 100644 index c9c8c894..00000000 --- a/net/src/main/java/com/zfoo/net/task/dispatcher/AbstractTaskDispatch.java +++ /dev/null @@ -1,37 +0,0 @@ -/* - * 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.task.dispatcher; - -import com.zfoo.protocol.util.StringUtils; - -/** - * @author jaysunxiao - * @version 3.0 - */ -public abstract class AbstractTaskDispatch implements ITaskDispatch { - - public static ITaskDispatch valueOf(String taskDispatchName) { - switch (taskDispatchName) { - case "random": - return new RandomTaskDispatch(); - case "sessionId": - return new SessionIdTaskDispatch(); - case "consistent-hash": - return new ConsistentHashTaskDispatch(); - default: - throw new RuntimeException(StringUtils.format("没有找到对应的taskDispatch[{}]", taskDispatchName)); - } - } - -} diff --git a/net/src/main/java/com/zfoo/net/task/dispatcher/ConsistentHashTaskDispatch.java b/net/src/main/java/com/zfoo/net/task/dispatcher/ConsistentHashTaskDispatch.java deleted file mode 100644 index 3e19e934..00000000 --- a/net/src/main/java/com/zfoo/net/task/dispatcher/ConsistentHashTaskDispatch.java +++ /dev/null @@ -1,55 +0,0 @@ -/* - * 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.task.dispatcher; - -import com.zfoo.net.session.model.AttributeType; -import com.zfoo.net.task.TaskBus; -import com.zfoo.net.task.model.PacketReceiverTask; -import com.zfoo.util.math.HashUtils; - -import java.util.concurrent.Executor; -import java.util.concurrent.ExecutorService; - -/** - * @author godotg - * @version 3.0 - */ -public class ConsistentHashTaskDispatch extends AbstractTaskDispatch { - - private static final ConsistentHashTaskDispatch INSTANCE = new ConsistentHashTaskDispatch(); - - public static ConsistentHashTaskDispatch getINSTANCE() { - return INSTANCE; - } - - @Override - public Executor getExecutor(ExecutorService[] executors, PacketReceiverTask packetReceiverTask) { - var attachment = packetReceiverTask.getAttachment(); - - if (attachment == null) { - var session = packetReceiverTask.getSession(); - var uid = session.getAttribute(AttributeType.UID); - - if (uid == null) { - return SessionIdTaskDispatch.getInstance().getExecutor(executors, packetReceiverTask); - } else { - return executors[TaskBus.executorIndex(HashUtils.fnvHash(uid))]; - } - } - - // 可见最终是根据附加包的信息选择服务端由哪个线程执行这个业务 - return executors[TaskBus.executorIndex(attachment.executorConsistentHash())]; - } - -} diff --git a/net/src/main/java/com/zfoo/net/task/dispatcher/ITaskDispatch.java b/net/src/main/java/com/zfoo/net/task/dispatcher/ITaskDispatch.java deleted file mode 100644 index 03be2830..00000000 --- a/net/src/main/java/com/zfoo/net/task/dispatcher/ITaskDispatch.java +++ /dev/null @@ -1,29 +0,0 @@ -/* - * 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.task.dispatcher; - -import com.zfoo.net.task.model.PacketReceiverTask; - -import java.util.concurrent.Executor; -import java.util.concurrent.ExecutorService; - -/** - * @author godotg - * @version 3.0 - */ -public interface ITaskDispatch { - - Executor getExecutor(ExecutorService[] executors, PacketReceiverTask packetReceiverTask); - -} diff --git a/net/src/main/java/com/zfoo/net/task/dispatcher/RandomTaskDispatch.java b/net/src/main/java/com/zfoo/net/task/dispatcher/RandomTaskDispatch.java deleted file mode 100644 index e9548f89..00000000 --- a/net/src/main/java/com/zfoo/net/task/dispatcher/RandomTaskDispatch.java +++ /dev/null @@ -1,40 +0,0 @@ -/* - * 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.task.dispatcher; - -import com.zfoo.net.task.TaskBus; -import com.zfoo.net.task.model.PacketReceiverTask; -import com.zfoo.util.math.RandomUtils; - -import java.util.concurrent.Executor; -import java.util.concurrent.ExecutorService; - -/** - * @author godotg - * @version 3.0 - */ -public class RandomTaskDispatch extends AbstractTaskDispatch { - - private static final RandomTaskDispatch INSTANCE = new RandomTaskDispatch(); - - public static RandomTaskDispatch getInstance() { - return INSTANCE; - } - - @Override - public Executor getExecutor(ExecutorService[] executors, PacketReceiverTask packetReceiverTask) { - return executors[TaskBus.executorIndex(RandomUtils.randomInt())]; - } - -} diff --git a/net/src/main/java/com/zfoo/net/task/dispatcher/SessionIdTaskDispatch.java b/net/src/main/java/com/zfoo/net/task/dispatcher/SessionIdTaskDispatch.java deleted file mode 100644 index 000637c4..00000000 --- a/net/src/main/java/com/zfoo/net/task/dispatcher/SessionIdTaskDispatch.java +++ /dev/null @@ -1,43 +0,0 @@ -/* - * 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.task.dispatcher; - -import com.zfoo.net.task.TaskBus; -import com.zfoo.net.task.model.PacketReceiverTask; -import com.zfoo.util.math.HashUtils; - -import java.util.concurrent.Executor; -import java.util.concurrent.ExecutorService; - -/** - * 同一个session总是分配到同一个线程池执行 - * - * @author godotg - * @version 3.0 - */ -public class SessionIdTaskDispatch extends AbstractTaskDispatch { - - private static final SessionIdTaskDispatch INSTANCE = new SessionIdTaskDispatch(); - - public static SessionIdTaskDispatch getInstance() { - return INSTANCE; - } - - @Override - public Executor getExecutor(ExecutorService[] executors, PacketReceiverTask packetReceiverTask) { - var session = packetReceiverTask.getSession(); - return executors[TaskBus.executorIndex(HashUtils.fnvHash(session.getSid()))]; - } - -} diff --git a/net/src/main/resources/net-1.0.xsd b/net/src/main/resources/net-1.0.xsd index 6393c409..6f80c14e 100644 --- a/net/src/main/resources/net-1.0.xsd +++ b/net/src/main/resources/net-1.0.xsd @@ -27,7 +27,6 @@ - diff --git a/net/src/test/resources/provider/provider_config.xml b/net/src/test/resources/provider/provider_config.xml index 62ca309b..e9fe9270 100644 --- a/net/src/test/resources/provider/provider_config.xml +++ b/net/src/test/resources/provider/provider_config.xml @@ -25,7 +25,7 @@ - +