mirror of
https://github.com/tiennm99/zfoo.git
synced 2026-09-03 06:19:37 +00:00
chore[rename]: rename RegisterVO to Register
This commit is contained in:
@@ -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() {
|
||||
|
||||
@@ -76,7 +76,7 @@ public class Consumer implements IConsumer {
|
||||
var protocolModule = ProtocolManager.moduleByProtocol(packet.getClass());
|
||||
var list = new ArrayList<Session>();
|
||||
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;
|
||||
|
||||
@@ -106,7 +106,7 @@ public class ConsistentHashLoadBalancer extends AbstractConsumerLoadBalancer {
|
||||
@Nullable
|
||||
private FastTreeMapIntLong updateModuleToConsistentHash(List<Session> 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();
|
||||
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -40,7 +40,7 @@ public interface IRegistry {
|
||||
|
||||
List<String> children(String path);
|
||||
|
||||
Set<RegisterVO> remoteProviderRegisterSet();
|
||||
Set<Register> remoteProviderRegisterSet();
|
||||
|
||||
/**
|
||||
* 监听path路径下的更新
|
||||
|
||||
+20
-20
@@ -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);
|
||||
}
|
||||
@@ -102,9 +102,9 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
*/
|
||||
private CuratorCache providerCuratorCache;
|
||||
/**
|
||||
* consumer需要消费的provider集合
|
||||
* 本地consumer需要消费的provider集合
|
||||
*/
|
||||
private final Set<RegisterVO> providerHashConsumerSet = new ConcurrentHashSet<>();
|
||||
private final Set<Register> 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<Session>();
|
||||
NetContext.getSessionManager().forEachClientSession(new Consumer<Session>() {
|
||||
@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<RegisterVO> remoteProviderRegisterSet() {
|
||||
public Set<Register> 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())) {
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user