This is an automated email from the ASF dual-hosted git repository.

Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git


The following commit(s) were added to refs/heads/master by this push:
     new f34847ff38 fix(admin): group URI registration by namespace (#7085)
f34847ff38 is described below

commit f34847ff382f7e67fd2060199c29977dca0ebbb5
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 09:57:26 2026 +0800

    fix(admin): group URI registration by namespace (#7085)
---
 .../subscriber/URIRegisterExecutorSubscriber.java  | 59 +++++++++-------------
 .../URIRegisterExecutorSubscriberTest.java         | 32 ++++++++++++
 2 files changed, 55 insertions(+), 36 deletions(-)

diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/disruptor/subscriber/URIRegisterExecutorSubscriber.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/disruptor/subscriber/URIRegisterExecutorSubscriber.java
index 254e268fa0..594743e9a9 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/disruptor/subscriber/URIRegisterExecutorSubscriber.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/disruptor/subscriber/URIRegisterExecutorSubscriber.java
@@ -28,10 +28,8 @@ import org.apache.shenyu.register.common.type.DataType;
 
 import java.util.Collection;
 import java.util.HashMap;
-import java.util.LinkedList;
 import java.util.List;
 import java.util.Map;
-import java.util.Objects;
 import java.util.Optional;
 import java.util.stream.Collectors;
 
@@ -70,43 +68,32 @@ public class URIRegisterExecutorSubscriber implements 
ExecutorTypeSubscriber<URI
                     .ifPresent(service -> {
                         final List<URIRegisterDTO> list = entry.getValue();
                         Map<String, List<URIRegisterDTO>> listMap = 
buildData(list);
-                        listMap.forEach((selectorName, uriList) -> {
-                            final List<URIRegisterDTO> register = new 
LinkedList<>();
-                            final List<URIRegisterDTO> heartbeat = new 
LinkedList<>();
-                            final List<URIRegisterDTO> offline = new 
LinkedList<>();
-                            for (URIRegisterDTO d : uriList) {
-                                final EventType eventType = d.getEventType();
-                                if (Objects.isNull(eventType) || 
EventType.REGISTER.equals(eventType)) {
-                                    // eventType is null, should be old 
versions
-                                    register.add(d);
-                                } else if 
(EventType.OFFLINE.equals(eventType)) {
-                                    offline.add(d);
-                                } else if 
(EventType.HEARTBEAT.equals(eventType)) {
-                                    heartbeat.add(d);
-                                }
-                            }
-                            if (CollectionUtils.isNotEmpty(register)) {
-                                
register.stream().map(URIRegisterDTO::getNamespaceId)
-                                        .filter(StringUtils::isNotBlank)
-                                        .findFirst()
-                                        .ifPresent(namespaceId -> 
service.registerURI(selectorName, register, namespaceId));
-                            }
-                            if (CollectionUtils.isNotEmpty(heartbeat)) {
-                                
heartbeat.stream().map(URIRegisterDTO::getNamespaceId)
-                                        .filter(StringUtils::isNotBlank)
-                                        .findFirst()
-                                        .ifPresent(namespaceId -> 
service.heartbeat(selectorName, heartbeat, namespaceId));
-                            }
-                            if (CollectionUtils.isNotEmpty(offline)) {
-                                
offline.stream().map(URIRegisterDTO::getNamespaceId)
-                                        .filter(StringUtils::isNotBlank)
-                                        .findFirst()
-                                        .ifPresent(namespaceId -> 
service.offline(selectorName, offline, namespaceId));
-                            }
-                        });
+                        listMap.forEach((selectorName, uriList) -> 
dispatchByNamespace(service, selectorName, uriList));
                     });
         }
     }
+
+    private void dispatchByNamespace(final ShenyuClientRegisterService 
service, final String selectorName, final List<URIRegisterDTO> uriList) {
+        Map<String, List<URIRegisterDTO>> groupByNamespace = uriList.stream()
+                .filter(dto -> StringUtils.isNotBlank(dto.getNamespaceId()))
+                
.collect(Collectors.groupingBy(URIRegisterDTO::getNamespaceId));
+        groupByNamespace.forEach((namespaceId, namespacedUris) -> {
+            Map<EventType, List<URIRegisterDTO>> groupByEventType = 
namespacedUris.stream()
+                    .collect(Collectors.groupingBy(dto -> 
Optional.ofNullable(dto.getEventType()).orElse(EventType.REGISTER)));
+            List<URIRegisterDTO> register = 
groupByEventType.get(EventType.REGISTER);
+            if (CollectionUtils.isNotEmpty(register)) {
+                service.registerURI(selectorName, register, namespaceId);
+            }
+            List<URIRegisterDTO> heartbeat = 
groupByEventType.get(EventType.HEARTBEAT);
+            if (CollectionUtils.isNotEmpty(heartbeat)) {
+                service.heartbeat(selectorName, heartbeat, namespaceId);
+            }
+            List<URIRegisterDTO> offline = 
groupByEventType.get(EventType.OFFLINE);
+            if (CollectionUtils.isNotEmpty(offline)) {
+                service.offline(selectorName, offline, namespaceId);
+            }
+        });
+    }
     
     private Map<String, List<URIRegisterDTO>> buildData(final 
Collection<URIRegisterDTO> dataList) {
         Map<String, List<URIRegisterDTO>> resultMap = new HashMap<>(8);
diff --git 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/disruptor/subscriber/URIRegisterExecutorSubscriberTest.java
 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/disruptor/subscriber/URIRegisterExecutorSubscriberTest.java
index 7dfd8462c5..ebecda9fbf 100644
--- 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/disruptor/subscriber/URIRegisterExecutorSubscriberTest.java
+++ 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/disruptor/subscriber/URIRegisterExecutorSubscriberTest.java
@@ -22,6 +22,7 @@ import org.apache.shenyu.common.constant.Constants;
 import org.apache.shenyu.common.enums.RpcTypeEnum;
 import org.apache.shenyu.common.exception.ShenyuException;
 import org.apache.shenyu.register.common.dto.URIRegisterDTO;
+import org.apache.shenyu.register.common.enums.EventType;
 import org.apache.shenyu.register.common.type.DataType;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
@@ -40,6 +41,8 @@ import java.util.Map;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.argThat;
+import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
@@ -74,6 +77,31 @@ public class URIRegisterExecutorSubscriberTest {
         uriRegisterExecutorSubscriber.executor(list);
         verify(service).registerURI(any(), any(), any());
     }
+
+    @Test
+    public void testExecutorGroupsUrisByNamespace() {
+        final String selectorName = "/test";
+        final String firstNamespace = "namespace-a";
+        final String secondNamespace = "namespace-b";
+        List<URIRegisterDTO> list = new ArrayList<>();
+        for (EventType eventType : EventType.values()) {
+            
list.add(URIRegisterDTO.builder().rpcType(RpcTypeEnum.HTTP.getName())
+                    
.contextPath(selectorName).namespaceId(firstNamespace).eventType(eventType).build());
+            
list.add(URIRegisterDTO.builder().rpcType(RpcTypeEnum.HTTP.getName())
+                    
.contextPath(selectorName).namespaceId(secondNamespace).eventType(eventType).build());
+        }
+        ShenyuClientRegisterService service = 
mock(ShenyuClientRegisterService.class);
+        
when(shenyuClientRegisterService.get(RpcTypeEnum.HTTP.getName())).thenReturn(service);
+
+        uriRegisterExecutorSubscriber.executor(list);
+
+        verify(service).registerURI(eq(selectorName), argThat(uris -> 
belongsToNamespace(uris, firstNamespace)), eq(firstNamespace));
+        verify(service).registerURI(eq(selectorName), argThat(uris -> 
belongsToNamespace(uris, secondNamespace)), eq(secondNamespace));
+        verify(service).heartbeat(eq(selectorName), argThat(uris -> 
belongsToNamespace(uris, firstNamespace)), eq(firstNamespace));
+        verify(service).heartbeat(eq(selectorName), argThat(uris -> 
belongsToNamespace(uris, secondNamespace)), eq(secondNamespace));
+        verify(service).offline(eq(selectorName), argThat(uris -> 
belongsToNamespace(uris, firstNamespace)), eq(firstNamespace));
+        verify(service).offline(eq(selectorName), argThat(uris -> 
belongsToNamespace(uris, secondNamespace)), eq(secondNamespace));
+    }
     
     @Test
     public void testBuildData() {
@@ -93,4 +121,8 @@ public class URIRegisterExecutorSubscriberTest {
             throw new ShenyuException(e.getCause());
         }
     }
+
+    private boolean belongsToNamespace(final List<URIRegisterDTO> uriList, 
final String namespaceId) {
+        return uriList.size() == 1 && 
namespaceId.equals(uriList.get(0).getNamespaceId());
+    }
 }

Reply via email to