feat[consumer]: write all the consumed servers to the data under the node

This commit is contained in:
godotg
2024-01-22 19:44:36 +08:00
parent efb479aebb
commit 3df1995f26
2 changed files with 50 additions and 31 deletions
@@ -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;
@@ -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<String>();
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)