diff --git a/net/src/main/java/com/zfoo/net/config/model/NetConfig.java b/net/src/main/java/com/zfoo/net/config/model/NetConfig.java index 81d1ba8f..254f1c35 100644 --- a/net/src/main/java/com/zfoo/net/config/model/NetConfig.java +++ b/net/src/main/java/com/zfoo/net/config/model/NetConfig.java @@ -13,7 +13,7 @@ package com.zfoo.net.config.model; -import com.zfoo.net.consumer.registry.RegisterVO; +import com.zfoo.net.consumer.registry.Register; import com.zfoo.protocol.generate.GenerateOperation; import java.util.Objects; @@ -61,8 +61,8 @@ public class NetConfig { private ConsumerConfig consumer; - public RegisterVO toLocalRegisterVO() { - return RegisterVO.valueOf(id, provider, consumer); + public Register toLocalRegister() { + return Register.valueOf(id, provider, consumer); } public String getId() { diff --git a/net/src/main/java/com/zfoo/net/consumer/Consumer.java b/net/src/main/java/com/zfoo/net/consumer/Consumer.java index 533516ce..2f3898f6 100644 --- a/net/src/main/java/com/zfoo/net/consumer/Consumer.java +++ b/net/src/main/java/com/zfoo/net/consumer/Consumer.java @@ -76,7 +76,7 @@ public class Consumer implements IConsumer { var protocolModule = ProtocolManager.moduleByProtocol(packet.getClass()); var list = new ArrayList(); NetContext.getSessionManager().forEachClientSession(session -> { - var consumerAttribute = session.getConsumerAttribute(); + var consumerAttribute = session.getConsumerRegister(); if (consumerAttribute == null) { return; } @@ -106,7 +106,7 @@ public class Consumer implements IConsumer { // 不同的服务提供者可能会提供同一个接口,消费者可能同时消费了这些提供了同一个接口的服务提供者,取第一个消费者的loadBalancer IConsumerLoadBalancer loadBalancer = null; for (var providerSession : providers) { - for (var provider : providerSession.getConsumerAttribute().getProviderConfig().getProviders()) { + for (var provider : providerSession.getConsumerRegister().getProviderConfig().getProviders()) { if (consumerLoadBalancerMap.containsKey(provider.getProvider())) { loadBalancer = consumerLoadBalancerMap.get(provider.getProvider()); break; diff --git a/net/src/main/java/com/zfoo/net/consumer/balancer/ConsistentHashLoadBalancer.java b/net/src/main/java/com/zfoo/net/consumer/balancer/ConsistentHashLoadBalancer.java index 009e6408..109260ba 100644 --- a/net/src/main/java/com/zfoo/net/consumer/balancer/ConsistentHashLoadBalancer.java +++ b/net/src/main/java/com/zfoo/net/consumer/balancer/ConsistentHashLoadBalancer.java @@ -106,7 +106,7 @@ public class ConsistentHashLoadBalancer extends AbstractConsumerLoadBalancer { @Nullable private FastTreeMapIntLong updateModuleToConsistentHash(List providers, ProtocolModule module) { var sessionStringList = providers.stream() - .map(session -> new Pair<>(session.getConsumerAttribute().toString(), session.getSid())) + .map(session -> new Pair<>(session.getConsumerRegister().toString(), session.getSid())) .sorted((a, b) -> a.getKey().compareTo(b.getKey())) .toList(); diff --git a/net/src/main/java/com/zfoo/net/consumer/event/ConsumerStartEvent.java b/net/src/main/java/com/zfoo/net/consumer/event/ConsumerStartEvent.java index c2fd4143..f2afc5a1 100644 --- a/net/src/main/java/com/zfoo/net/consumer/event/ConsumerStartEvent.java +++ b/net/src/main/java/com/zfoo/net/consumer/event/ConsumerStartEvent.java @@ -14,7 +14,7 @@ package com.zfoo.net.consumer.event; import com.zfoo.event.model.IEvent; -import com.zfoo.net.consumer.registry.RegisterVO; +import com.zfoo.net.consumer.registry.Register; import com.zfoo.net.session.Session; /** @@ -22,22 +22,22 @@ import com.zfoo.net.session.Session; */ public class ConsumerStartEvent implements IEvent { - private RegisterVO providerRegisterVO; + private Register providerRegister; private Session session; - public static ConsumerStartEvent valueOf(RegisterVO providerRegisterVO, Session session) { + public static ConsumerStartEvent valueOf(Register providerRegister, Session session) { var event = new ConsumerStartEvent(); - event.providerRegisterVO = providerRegisterVO; + event.providerRegister = providerRegister; event.session = session; return event; } - public RegisterVO getProviderRegisterVO() { - return providerRegisterVO; + public Register getProviderRegister() { + return providerRegister; } - public void setProviderRegisterVO(RegisterVO providerRegisterVO) { - this.providerRegisterVO = providerRegisterVO; + public void setProviderRegister(Register providerRegister) { + this.providerRegister = providerRegister; } public Session getSession() { diff --git a/net/src/main/java/com/zfoo/net/consumer/registry/IRegistry.java b/net/src/main/java/com/zfoo/net/consumer/registry/IRegistry.java index 59259649..f087b996 100644 --- a/net/src/main/java/com/zfoo/net/consumer/registry/IRegistry.java +++ b/net/src/main/java/com/zfoo/net/consumer/registry/IRegistry.java @@ -40,7 +40,7 @@ public interface IRegistry { List children(String path); - Set remoteProviderRegisterSet(); + Set remoteProviderRegisterSet(); /** * 监听path路径下的更新 diff --git a/net/src/main/java/com/zfoo/net/consumer/registry/RegisterVO.java b/net/src/main/java/com/zfoo/net/consumer/registry/Register.java similarity index 81% rename from net/src/main/java/com/zfoo/net/consumer/registry/RegisterVO.java rename to net/src/main/java/com/zfoo/net/consumer/registry/Register.java index 8205692d..f628ee3f 100644 --- a/net/src/main/java/com/zfoo/net/consumer/registry/RegisterVO.java +++ b/net/src/main/java/com/zfoo/net/consumer/registry/Register.java @@ -32,9 +32,9 @@ import java.util.Objects; /** * @author godotg */ -public class RegisterVO { +public class Register { - private static final Logger logger = LoggerFactory.getLogger(RegisterVO.class); + private static final Logger logger = LoggerFactory.getLogger(Register.class); private static final String LOCAL_UUID = UuidUtils.getUUID(); @@ -45,34 +45,34 @@ public class RegisterVO { // 服务消费者配置 private ConsumerConfig consumerConfig; - public static boolean providerHasConsumer(RegisterVO providerVO, RegisterVO consumerVO) { - if (Objects.isNull(providerVO) || Objects.isNull(providerVO.providerConfig) || CollectionUtils.isEmpty(providerVO.providerConfig.getProviders()) - || Objects.isNull(consumerVO) || Objects.isNull(consumerVO.consumerConfig) || CollectionUtils.isEmpty(consumerVO.consumerConfig.getConsumers())) { + public static boolean providerHasConsumer(Register providerRegister, Register consumerRegister) { + if (Objects.isNull(providerRegister) || Objects.isNull(providerRegister.providerConfig) || CollectionUtils.isEmpty(providerRegister.providerConfig.getProviders()) + || Objects.isNull(consumerRegister) || Objects.isNull(consumerRegister.consumerConfig) || CollectionUtils.isEmpty(consumerRegister.consumerConfig.getConsumers())) { return false; } - for (var provider : providerVO.getProviderConfig().getProviders()) { - if (consumerVO.getConsumerConfig().getConsumers().stream().anyMatch(it -> it.getConsumer().equals(provider.getProvider()))) { + for (var provider : providerRegister.getProviderConfig().getProviders()) { + if (consumerRegister.getConsumerConfig().getConsumers().stream().anyMatch(it -> it.getConsumer().equals(provider.getProvider()))) { return true; } } return false; } - public static RegisterVO valueOf(String id, ProviderConfig providerConfig, ConsumerConfig consumerConfig) { - RegisterVO config = new RegisterVO(); - config.id = id; - config.providerConfig = providerConfig; - config.consumerConfig = consumerConfig; - return config; + public static Register valueOf(String id, ProviderConfig providerConfig, ConsumerConfig consumerConfig) { + Register register = new Register(); + register.id = id; + register.providerConfig = providerConfig; + register.consumerConfig = consumerConfig; + return register; } @Nullable - public static RegisterVO parseString(String str) { + public static Register parseString(String str) { try { - var vo = new RegisterVO(); + var register = new Register(); var splits = str.split("\\|"); - vo.id = splits[0].trim(); + register.id = splits[0].trim(); String providerAddress = null; @@ -80,16 +80,16 @@ public class RegisterVO { var s = splits[i].trim(); if (s.startsWith("provider")) { var providerModules = parseProviderModules(s); - vo.providerConfig = ProviderConfig.valueOf(providerAddress, providerModules); + register.providerConfig = ProviderConfig.valueOf(providerAddress, providerModules); } else if (s.startsWith("consumer")) { var consumerModules = parseConsumerModules(s); - vo.consumerConfig = ConsumerConfig.valueOf(consumerModules); + register.consumerConfig = ConsumerConfig.valueOf(consumerModules); } else { providerAddress = s; } } - return vo; + return register; } catch (Exception e) { logger.error(ExceptionUtils.getMessage(e)); return null; @@ -194,7 +194,7 @@ public class RegisterVO { if (o == null || getClass() != o.getClass()) { return false; } - RegisterVO that = (RegisterVO) o; + Register that = (Register) o; return Objects.equals(id, that.id) && Objects.equals(providerConfig, that.providerConfig) && Objects.equals(consumerConfig, that.consumerConfig); } 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 1a84b495..bfaad279 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 @@ -102,9 +102,9 @@ public class ZookeeperRegistry implements IRegistry { */ private CuratorCache providerCuratorCache; /** - * consumer需要消费的provider集合 + * 本地consumer需要消费的provider集合 */ - private final Set providerHashConsumerSet = new ConcurrentHashSet<>(); + private final Set providerRegisterSet = new ConcurrentHashSet<>(); /** * addListener中的cache全部会被添加到这个集合中,这个集合不包括providerCuratorCache */ @@ -308,21 +308,21 @@ public class ZookeeperRegistry implements IRegistry { break; case NODE_CREATED: // 意味着有可能来了自己作为消费者需要关心的服务提供者 var providerStr = StringUtils.substringAfterFirst(newData.getPath(), PROVIDER_ROOT_PATH + StringUtils.SLASH); - var provider = RegisterVO.parseString(providerStr); - var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegisterVO(); + var providerRegister = Register.parseString(providerStr); + var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegister(); // 如果启动的Consumer是自己关心的Consumer,那么就会接下来尝试连接他们 // 这意味着:如果有多个Consumer启动,那么最后将全部连接上去 - if (RegisterVO.providerHasConsumer(provider, localRegisterVO)) { - providerHashConsumerSet.add(provider); + if (Register.providerHasConsumer(providerRegister, localRegisterVO)) { + providerRegisterSet.add(providerRegister); checkConsumer(); logger.info("Discover new subscription service of provider [{}]", providerStr); } break; case NODE_DELETED: var oldProviderStr = StringUtils.substringAfterFirst(oldData.getPath(), PROVIDER_ROOT_PATH + StringUtils.SLASH); - var oldProvider = RegisterVO.parseString(oldProviderStr); - if (providerHashConsumerSet.contains(oldProvider)) { - providerHashConsumerSet.remove(oldProvider); + var oldProvider = Register.parseString(oldProviderStr); + if (providerRegisterSet.contains(oldProvider)) { + providerRegisterSet.remove(oldProvider); checkConsumer(); logger.info("Unsubscribe from the service of provider [{}]", oldProviderStr); } @@ -361,7 +361,7 @@ public class ZookeeperRegistry implements IRegistry { * 如果自己是服务提供者,就把自己注册上去 */ private void initLocalProvider() throws Exception { - var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegisterVO(); + var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegister(); if (Objects.nonNull(localRegisterVO.getProviderConfig())) { var localProviderVoStr = localRegisterVO.toProviderString(); var localProviderPath = PROVIDER_ROOT_PATH + StringUtils.SLASH + localProviderVoStr; @@ -399,21 +399,21 @@ public class ZookeeperRegistry implements IRegistry { */ private void initConsumerCache() throws Exception { // /zfoo/provider/applicationNameTest | 192.168.1.104:12400 | provider:[myProviderModule-provider1, myProviderModule-provider2] - var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegisterVO(); + var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegister(); // 初始化providerCacheSet // 遍历provider下注册的所有节点 var remoteProviderSet = curator.getChildren().forPath(PROVIDER_ROOT_PATH).stream() .filter(it -> StringUtils.isNotBlank(it) && !"null".equals(it)) - .map(it -> RegisterVO.parseString(it)) + .map(it -> Register.parseString(it)) .filter(it -> Objects.nonNull(it)) // 检查是否这个节点是自己关心的节点 - .filter(it -> RegisterVO.providerHasConsumer(it, localRegisterVO)) + .filter(it -> Register.providerHasConsumer(it, localRegisterVO)) .collect(Collectors.toSet()); - providerHashConsumerSet.clear(); + providerRegisterSet.clear(); // 将自己关心的节点存起来,接下来,将会开启TcpClient去连接这些Provider,连接上后,将会把这个session保存到ClientSessionMap中 - providerHashConsumerSet.addAll(remoteProviderSet); + providerRegisterSet.addAll(remoteProviderSet); // 初始化consumer,providerCacheSet改变会导致消费者改变 // 如果自己没有连接上远程消费者,则会一直尝试连接 @@ -445,17 +445,17 @@ public class ZookeeperRegistry implements IRegistry { return; } - logger.info("start using providerHashConsumerSet:{} to check [consumer:{}]", providerHashConsumerSet, NetContext.getSessionManager().clientSessionSize()); + logger.info("start using providerHashConsumerSet:{} to check [consumer:{}]", providerRegisterSet, NetContext.getSessionManager().clientSessionSize()); var recheckFlag = false; - for (var providerCache : providerHashConsumerSet) { + for (var providerCache : providerRegisterSet) { // 先排除已经启动的consumer var consumerClientList = new ArrayList(); NetContext.getSessionManager().forEachClientSession(new Consumer() { @Override public void accept(Session session) { - if (session.getConsumerAttribute() != null && session.getConsumerAttribute().equals(providerCache)) { + if (session.getConsumerRegister() != null && session.getConsumerRegister().equals(providerCache)) { consumerClientList.add(session); } } @@ -486,7 +486,7 @@ public class ZookeeperRegistry implements IRegistry { continue; } - session.setConsumerAttribute(providerCache); + session.setConsumerRegister(providerCache); EventBus.post(ConsumerStartEvent.valueOf(providerCache, session)); } catch (Throwable t) { logger.error("[consumer:{}] failed to start, wait [{}] seconds to recheck consumer", providerCache, RETRY_SECONDS, t); @@ -497,7 +497,7 @@ public class ZookeeperRegistry implements IRegistry { // 将自己的消费者消息写到 /consumer 的临时节点下 var consumerConfig = NetContext.getConfigManager().getLocalConfig().getConsumer(); if (consumerConfig != null && CollectionUtils.isNotEmpty(consumerConfig.getConsumers())) { - var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegisterVO(); + var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegister(); var path = CONSUMER_ROOT_PATH + StringUtils.SLASH + localRegisterVO.toConsumerString(); try { var stat = curator.checkExists().forPath(path); @@ -617,11 +617,11 @@ public class ZookeeperRegistry implements IRegistry { * @return */ @Override - public Set remoteProviderRegisterSet() { + public Set remoteProviderRegisterSet() { try { var remoteProviderSet = curator.getChildren().forPath(PROVIDER_ROOT_PATH).stream() .filter(it -> StringUtils.isNotBlank(it) && !"null".equals(it)) - .map(it -> RegisterVO.parseString(it)) + .map(it -> Register.parseString(it)) .filter(it -> Objects.nonNull(it)) .collect(Collectors.toSet()); return remoteProviderSet; @@ -691,7 +691,7 @@ public class ZookeeperRegistry implements IRegistry { return; } try { - var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegisterVO(); + var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegister(); if (curator.getState() == CuratorFrameworkState.STARTED) { // 删除服务提供者的临时节点 if (Objects.nonNull(localRegisterVO.getProviderConfig())) { 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 de68f220..74a8cce5 100644 --- a/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java +++ b/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java @@ -54,7 +54,7 @@ public class ClientRouteHandler extends BaseRouteHandler { EventBus.post(ClientSessionInactiveEvent.valueOf(session)); // 如果是消费者inactive,还需要触发客户端消费者检查事件,以便重新连接 - if (session.getConsumerAttribute() != null) { + if (session.getConsumerRegister() != null) { NetContext.getConfigManager().getRegistry().checkConsumer(); } diff --git a/net/src/main/java/com/zfoo/net/session/Session.java b/net/src/main/java/com/zfoo/net/session/Session.java index 433e1394..1422f5b1 100644 --- a/net/src/main/java/com/zfoo/net/session/Session.java +++ b/net/src/main/java/com/zfoo/net/session/Session.java @@ -12,7 +12,7 @@ package com.zfoo.net.session; -import com.zfoo.net.consumer.registry.RegisterVO; +import com.zfoo.net.consumer.registry.Register; import com.zfoo.protocol.util.StringUtils; import io.netty.channel.Channel; @@ -45,7 +45,7 @@ public class Session implements Closeable { * EN:Session extra parameters * CN:Session附带的属性参数,消费者的属性 */ - private RegisterVO consumerAttribute = null; + private Register consumerRegister = null; public Session(Channel channel) { if (channel == null) { @@ -98,11 +98,11 @@ public class Session implements Closeable { this.uid = uid; } - public RegisterVO getConsumerAttribute() { - return consumerAttribute; + public Register getConsumerRegister() { + return consumerRegister; } - public void setConsumerAttribute(RegisterVO consumerAttribute) { - this.consumerAttribute = consumerAttribute; + public void setConsumerRegister(Register consumerRegister) { + this.consumerRegister = consumerRegister; } } diff --git a/net/src/test/java/com/zfoo/net/config/RegistryTest.java b/net/src/test/java/com/zfoo/net/config/RegistryTest.java index b9785206..1f673f80 100644 --- a/net/src/test/java/com/zfoo/net/config/RegistryTest.java +++ b/net/src/test/java/com/zfoo/net/config/RegistryTest.java @@ -17,9 +17,8 @@ import com.zfoo.net.config.model.ConsumerConfig; import com.zfoo.net.config.model.ConsumerModule; import com.zfoo.net.config.model.ProviderConfig; import com.zfoo.net.config.model.ProviderModule; -import com.zfoo.net.consumer.registry.RegisterVO; +import com.zfoo.net.consumer.registry.Register; import com.zfoo.net.core.HostAndPort; -import com.zfoo.protocol.registration.ProtocolModule; import io.netty.util.NetUtil; import org.junit.Assert; import org.junit.Test; @@ -48,14 +47,14 @@ public class RegistryTest { // 服务消费者配置:这个是没Ip的 var consumerConfig = ConsumerConfig.valueOf(consumerModules); - var vo = RegisterVO.valueOf("test", providerConfig, consumerConfig); - var voStr = vo.toString(); + var register = Register.valueOf("test", providerConfig, consumerConfig); + var voStr = register.toString(); // test | 127.0.0.1:80 | provider:[100-aaa-a, 120-bbb-b] | consumer:[100-aaa-random-a, 120-bbb-random-b] System.out.println(voStr); - var newVo = RegisterVO.parseString(voStr); - Assert.assertEquals(vo, newVo); + var newRegister = Register.parseString(voStr); + Assert.assertEquals(register, newRegister); // /127.0.0.1 System.out.println(NetUtil.LOCALHOST); diff --git a/net/src/test/java/com/zfoo/net/core/gateway/GatewayProviderController.java b/net/src/test/java/com/zfoo/net/core/gateway/GatewayProviderController.java index 56d52185..86953310 100644 --- a/net/src/test/java/com/zfoo/net/core/gateway/GatewayProviderController.java +++ b/net/src/test/java/com/zfoo/net/core/gateway/GatewayProviderController.java @@ -40,7 +40,7 @@ public class GatewayProviderController { logger.info("provider receive [packet:{}] from client", JsonUtils.object2String(request)); var response = new GatewayToProviderResponse(); - response.setMessage(StringUtils.format("Hello, this is the [provider:{}] response!", NetContext.getConfigManager().getLocalConfig().toLocalRegisterVO().toString())); + response.setMessage(StringUtils.format("Hello, this is the [provider:{}] response!", NetContext.getConfigManager().getLocalConfig().toLocalRegister().toString())); NetContext.getRouter().send(session, response, gatewayAttachment); } diff --git a/net/src/test/java/com/zfoo/net/core/provider/ProviderController.java b/net/src/test/java/com/zfoo/net/core/provider/ProviderController.java index b8a49d9c..7db3fd7c 100644 --- a/net/src/test/java/com/zfoo/net/core/provider/ProviderController.java +++ b/net/src/test/java/com/zfoo/net/core/provider/ProviderController.java @@ -36,7 +36,7 @@ public class ProviderController { logger.info("provider receive [packet:{}] from consumer", JsonUtils.object2String(ask)); var response = new ProviderMessAnswer(); - response.setMessage(StringUtils.format("Hello, this is the [provider:{}] answer!", NetContext.getConfigManager().getLocalConfig().toLocalRegisterVO().toString())); + response.setMessage(StringUtils.format("Hello, this is the [provider:{}] answer!", NetContext.getConfigManager().getLocalConfig().toLocalRegister().toString())); NetContext.getRouter().send(session, response); }