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());
+ }
}