This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new d10860dce8 [INLONG-8330][Manager] Fix the information returned by the
getAllConfig interface is incorrect (#8331)
d10860dce8 is described below
commit d10860dce83c7b2e2853fc74c6d77495ba1af471
Author: fuweng11 <[email protected]>
AuthorDate: Wed Jun 28 11:39:06 2023 +0800
[INLONG-8330][Manager] Fix the information returned by the getAllConfig
interface is incorrect (#8331)
---
.../service/repository/DataProxyConfigRepository.java | 12 +++++++++---
.../service/repository/DataProxyConfigRepositoryV2.java | 12 +++++++++---
2 files changed, 18 insertions(+), 6 deletions(-)
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
index 263b91c52e..b5dc9fabf0 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
@@ -48,6 +48,7 @@ import org.apache.inlong.manager.pojo.sink.SinkPageRequest;
import org.apache.inlong.manager.service.core.SortConfigLoader;
import com.google.common.base.Splitter;
+import com.google.common.collect.Sets;
import com.google.gson.Gson;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
@@ -71,6 +72,7 @@ import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Optional;
+import java.util.Set;
import java.util.Timer;
import java.util.TimerTask;
import java.util.concurrent.ConcurrentHashMap;
@@ -255,9 +257,13 @@ public class DataProxyConfigRepository implements
IRepository {
}
Map<String, String> tagMap =
MAP_SPLITTER.split(cacheCluster.getExtTag());
String producerTag =
tagMap.getOrDefault(CACHE_CLUSTER_PRODUCER_TAG, Boolean.TRUE.toString());
- if (StringUtils.equalsIgnoreCase(producerTag,
Boolean.TRUE.toString())) {
- cacheClusterMap.computeIfAbsent(cacheCluster.getClusterTags(),
k -> new HashMap<>())
- .computeIfAbsent(cacheCluster.getExtTag(), k -> new
ArrayList<>()).add(cacheCluster);
+ if (StringUtils.equalsIgnoreCase(producerTag,
Boolean.TRUE.toString()) && StringUtils.isNotBlank(
+ cacheCluster.getClusterTags())) {
+ Set<String> clusterTags =
Sets.newHashSet(cacheCluster.getClusterTags().split(InlongConstants.COMMA));
+ clusterTags.forEach(clusterTag -> {
+ cacheClusterMap.computeIfAbsent(clusterTag, k -> new
HashMap<>())
+ .computeIfAbsent(cacheCluster.getExtTag(), k ->
new ArrayList<>()).add(cacheCluster);
+ });
}
}
// mark cache cluster to proxy cluster
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepositoryV2.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepositoryV2.java
index a49dbdd4b1..d501e62f96 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepositoryV2.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepositoryV2.java
@@ -48,6 +48,7 @@ import org.apache.inlong.manager.pojo.sink.SinkPageRequest;
import org.apache.inlong.manager.service.core.SortConfigLoader;
import com.google.common.base.Splitter;
+import com.google.common.collect.Sets;
import com.google.gson.Gson;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
@@ -71,6 +72,7 @@ import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Optional;
+import java.util.Set;
import java.util.Timer;
import java.util.TimerTask;
import java.util.concurrent.ConcurrentHashMap;
@@ -252,9 +254,13 @@ public class DataProxyConfigRepositoryV2 implements
IRepository {
}
Map<String, String> tagMap =
MAP_SPLITTER.split(cacheCluster.getExtTag());
String producerTag =
tagMap.getOrDefault(CACHE_CLUSTER_PRODUCER_TAG, Boolean.TRUE.toString());
- if (StringUtils.equalsIgnoreCase(producerTag,
Boolean.TRUE.toString())) {
- cacheClusterMap.computeIfAbsent(cacheCluster.getClusterTags(),
k -> new HashMap<>())
- .computeIfAbsent(cacheCluster.getExtTag(), k -> new
ArrayList<>()).add(cacheCluster);
+ if (StringUtils.equalsIgnoreCase(producerTag,
Boolean.TRUE.toString()) && StringUtils.isNotBlank(
+ cacheCluster.getClusterTags())) {
+ Set<String> clusterTags =
Sets.newHashSet(cacheCluster.getClusterTags().split(InlongConstants.COMMA));
+ clusterTags.forEach(clusterTag -> {
+ cacheClusterMap.computeIfAbsent(clusterTag, k -> new
HashMap<>())
+ .computeIfAbsent(cacheCluster.getExtTag(), k ->
new ArrayList<>()).add(cacheCluster);
+ });
}
}
// mark cache cluster to proxy cluster