From 3df1995f267378b4c396a35faa34a9da623ab627 Mon Sep 17 00:00:00 2001 From: godotg Date: Mon, 22 Jan 2024 19:44:36 +0800 Subject: [PATCH] feat[consumer]: write all the consumed servers to the data under the node --- .../zfoo/net/consumer/registry/Register.java | 37 ++++++++++++---- .../consumer/registry/ZookeeperRegistry.java | 44 +++++++++---------- 2 files changed, 50 insertions(+), 31 deletions(-) diff --git a/net/src/main/java/com/zfoo/net/consumer/registry/Register.java b/net/src/main/java/com/zfoo/net/consumer/registry/Register.java index ec80a784..375e7c19 100644 --- a/net/src/main/java/com/zfoo/net/consumer/registry/Register.java +++ b/net/src/main/java/com/zfoo/net/consumer/registry/Register.java @@ -123,15 +123,6 @@ public class Register { } public String toProviderString() { - return toString(); - } - - public String toConsumerString() { - return this + StringUtils.SPACE + StringUtils.VERTICAL_BAR + StringUtils.SPACE + LOCAL_UUID; - } - - @Override - public String toString() { var builder = new StringBuilder(); // 模块模块名 @@ -173,6 +164,34 @@ public class Register { return builder.toString(); } + public String toProviderSimple() { + var builder = new StringBuilder(); + + // 模块模块名 + builder.append(id); + + // 服务提供者相关配置信息 + if (Objects.nonNull(providerConfig)) { + var providerAddress = providerConfig.getAddress(); + if (StringUtils.isBlank(providerAddress)) { + throw new RuntimeException(StringUtils.format("The address of provider Config cannot be empty")); + } + builder.append(StringUtils.SPACE).append(StringUtils.VERTICAL_BAR).append(StringUtils.SPACE); + // 服务提供者地址 + builder.append(providerAddress); + } + return builder.toString(); + } + + public String toConsumerString() { + return this + StringUtils.SPACE + StringUtils.VERTICAL_BAR + StringUtils.SPACE + LOCAL_UUID; + } + + @Override + public String toString() { + return toProviderString(); + } + public String getId() { return id; 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 908adf92..4b01aa4a 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 @@ -490,36 +490,36 @@ public class ZookeeperRegistry implements IRegistry { session.setConsumerRegister(providerCache); logger.info("Consumer starts consuming the provider:[{}]", providerCache); EventBus.post(ConsumerStartEvent.valueOf(providerCache, session)); + // 将自己的消费者消息写到 /consumer 的临时节点下 + updateConsumerData(); } catch (Throwable t) { logger.error("[consumer:{}] failed to start, wait [{}] seconds to recheck consumer", providerCache, RETRY_SECONDS, t); recheckFlag = true; } } - // 将自己的消费者消息写到 /consumer 的临时节点下 - var consumerConfig = NetContext.getConfigManager().getLocalConfig().getConsumer(); - if (consumerConfig != null && CollectionUtils.isNotEmpty(consumerConfig.getConsumers())) { - var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegister(); - var path = CONSUMER_ROOT_PATH + StringUtils.SLASH + localRegisterVO.toConsumerString(); - try { - var stat = curator.checkExists().forPath(path); - if (Objects.isNull(stat)) { - curator.create() - .withMode(CreateMode.EPHEMERAL) - .forPath(path); - } else { - curator.setData().forPath(path); - } - } catch (Exception e) { - logger.error("consumer:[{}] writing to Zookeeper failed", path, e); - recheckFlag = true; - } - } - if (recheckFlag) { SchedulerBus.schedule(() -> checkConsumer(), RETRY_SECONDS, TimeUnit.SECONDS); } + } + private void updateConsumerData() { + // 将自己的消费者消息写到 /consumer 的临时节点下 + var localRegisterVO = NetContext.getConfigManager().getLocalConfig().toLocalRegister(); + var path = CONSUMER_ROOT_PATH + StringUtils.SLASH + localRegisterVO.toConsumerString(); + var list = new ArrayList(); + NetContext.getSessionManager().forEachClientSession(session -> { + var consumerAttribute = session.getConsumerRegister(); + if (consumerAttribute == null) { + return; + } + var providerConfig = consumerAttribute.getProviderConfig(); + if (providerConfig == null) { + return; + } + list.add(consumerAttribute.toProviderSimple()); + }); + addData(path, StringUtils.bytes(JsonUtils.object2String(list)), CreateMode.EPHEMERAL); } /** @@ -532,9 +532,9 @@ public class ZookeeperRegistry implements IRegistry { @Override public void addData(String path, byte[] bytes, CreateMode mode) { try { - var providerStat = curator.checkExists().forPath(path); + var stat = curator.checkExists().forPath(path); - if (Objects.isNull(providerStat)) { + if (Objects.isNull(stat)) { curator.create() .creatingParentsIfNeeded() .withMode(mode)