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 5d599ba2..78c81399 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 @@ -493,7 +493,7 @@ public class ZookeeperRegistry implements IRegistry { } else { // 连接上了服务提供者 session.putAttribute(AttributeType.CONSUMER, providerCache); - EventBus.asyncSubmit(ConsumerStartEvent.valueOf(providerCache, session)); + EventBus.submit(ConsumerStartEvent.valueOf(providerCache, session)); try { var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegisterVO(); diff --git a/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java b/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java index 89438e2f..e15d7c93 100644 --- a/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java +++ b/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java @@ -50,7 +50,7 @@ public class ClientRouteHandler extends BaseRouteHandler { var consumeAttribute = session.getAttribute(AttributeType.CONSUMER); NetContext.getSessionManager().removeClientSession(session); - EventBus.asyncSubmit(ClientSessionInactiveEvent.valueOf(session)); + EventBus.submit(ClientSessionInactiveEvent.valueOf(session)); // 如果是消费者inactive,还需要触发客户端消费者检查事件,以便重新连接 if (consumeAttribute != null) { diff --git a/net/src/main/java/com/zfoo/net/handler/GatewayRouteHandler.java b/net/src/main/java/com/zfoo/net/handler/GatewayRouteHandler.java index 23822451..6a9e8301 100644 --- a/net/src/main/java/com/zfoo/net/handler/GatewayRouteHandler.java +++ b/net/src/main/java/com/zfoo/net/handler/GatewayRouteHandler.java @@ -130,7 +130,7 @@ public class GatewayRouteHandler extends ServerRouteHandler { var uid = (Long) session.getAttribute(AttributeType.UID); // 连接到网关的客户端断开了连接 - EventBus.asyncSubmit(GatewaySessionInactiveEvent.valueOf(sid, uid == null ? 0 : uid.longValue())); + EventBus.submit(GatewaySessionInactiveEvent.valueOf(sid, uid == null ? 0 : uid.longValue())); super.channelInactive(ctx); } diff --git a/net/src/main/java/com/zfoo/net/handler/ServerRouteHandler.java b/net/src/main/java/com/zfoo/net/handler/ServerRouteHandler.java index 18b40d22..00697b1c 100644 --- a/net/src/main/java/com/zfoo/net/handler/ServerRouteHandler.java +++ b/net/src/main/java/com/zfoo/net/handler/ServerRouteHandler.java @@ -48,7 +48,7 @@ public class ServerRouteHandler extends BaseRouteHandler { return; } NetContext.getSessionManager().removeServerSession(session); - EventBus.asyncSubmit(ServerSessionInactiveEvent.valueOf(session)); + EventBus.submit(ServerSessionInactiveEvent.valueOf(session)); logger.warn("server channel is inactive {}", SessionUtils.sessionSimpleInfo(ctx)); } } 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 e979ca79..e1035a87 100644 --- a/net/src/main/java/com/zfoo/net/router/Router.java +++ b/net/src/main/java/com/zfoo/net/router/Router.java @@ -135,7 +135,7 @@ public class Router implements IRouter { return; } gatewaySession.putAttribute(AttributeType.UID, uid); - EventBus.asyncSubmit(AuthUidToGatewayEvent.valueOf(gatewaySession.getSid(), uid)); + EventBus.submit(AuthUidToGatewayEvent.valueOf(gatewaySession.getSid(), uid)); NetContext.getRouter().send(session, AuthUidToGatewayConfirm.valueOf(uid), new GatewayAttachment(gatewaySession, null)); return; @@ -346,7 +346,7 @@ public class Router implements IRouter { // 这个在哪个线程处理取决于:这个上层的PacketReceiverTask被丢到了哪个线程中 PacketBus.submit(session, packet, attachment); } catch (Exception e) { - EventBus.syncSubmit(ServerExceptionEvent.valueOf(session, packet, attachment, e)); + EventBus.submit(ServerExceptionEvent.valueOf(session, packet, attachment, e)); logger.error(StringUtils.format("e[uid:{}][sid:{}]未知exception异常", session.getAttribute(AttributeType.UID), session.getSid(), e.getMessage()), e); } catch (Throwable t) { logger.error(StringUtils.format("e[uid:{}][sid:{}]未知error错误", session.getAttribute(AttributeType.UID), session.getSid(), t.getMessage()), t);