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

Reply via email to