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 6b85c6ddd6 fix(dubbo): offload reference initialization from request
threads (#7273)
6b85c6ddd6 is described below
commit 6b85c6ddd6becf1262f0932bbe7a38bbf799fa2d
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 06:27:58 2026 +0800
fix(dubbo): offload reference initialization from request threads (#7273)
---
.../dubbo/proxy/ApacheDubboProxyService.java | 19 ++++++
.../dubbo/proxy/ApacheDubboProxyServiceTest.java | 75 ++++++++++++++++++++++
2 files changed, 94 insertions(+)
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyService.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyService.java
index 3dd5401f48..2fcab86f5b 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyService.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyService.java
@@ -43,9 +43,12 @@ import
org.apache.shenyu.plugin.dubbo.common.param.DubboParamResolveService;
import org.springframework.util.ObjectUtils;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Schedulers;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
@@ -79,6 +82,22 @@ public class ApacheDubboProxyService {
* @throws ShenyuException the shenyu exception
*/
public Mono<Object> genericInvoker(final String body, final MetaData
metaData, final SelectorData selectorData, final RuleData ruleData, final
ServerWebExchange exchange) throws ShenyuException {
+ Map<String, Object> attachments = new
HashMap<>(RpcContext.getClientAttachment().getObjectAttachments());
+ return Mono.defer(() -> {
+ Map<String, Object> previous = new
HashMap<>(RpcContext.getClientAttachment().getObjectAttachments());
+ try {
+ RpcContext.getClientAttachment().setObjectAttachments(new
HashMap<>(attachments));
+ return invokeOnWorker(body, metaData, selectorData, ruleData,
exchange);
+ } finally {
+ // Invocation and future lookup are synchronous; do not retain
request data on a pooled worker.
+ RpcContext.getClientAttachment().clearAttachments();
+
RpcContext.getClientAttachment().setObjectAttachments(previous);
+ }
+ })
+ .subscribeOn(Schedulers.boundedElastic());
+ }
+
+ private Mono<Object> invokeOnWorker(final String body, final MetaData
metaData, final SelectorData selectorData, final RuleData ruleData, final
ServerWebExchange exchange) {
ReferenceConfig<GenericService> reference =
this.getReferenceConfig(selectorData, ruleData, metaData, exchange);
GenericService genericService = reference.get();
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyServiceTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyServiceTest.java
index f64c3d3408..cea85823ab 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyServiceTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyServiceTest.java
@@ -21,8 +21,11 @@ import com.google.common.cache.LoadingCache;
import org.apache.commons.lang3.tuple.ImmutablePair;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.dubbo.config.ReferenceConfig;
+import org.apache.dubbo.rpc.RpcContext;
+import org.apache.dubbo.rpc.RpcContextAttachment;
import org.apache.dubbo.rpc.service.GenericService;
import org.apache.shenyu.common.dto.MetaData;
+import org.apache.shenyu.common.constant.Constants;
import org.apache.shenyu.common.dto.RuleData;
import org.apache.shenyu.common.dto.SelectorData;
import org.apache.shenyu.common.enums.RpcTypeEnum;
@@ -42,9 +45,18 @@ import org.springframework.web.server.ServerWebExchange;
import java.lang.reflect.Field;
import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.atomic.AtomicReference;
+import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Schedulers;
+import reactor.test.StepVerifier;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* The Test Case For ApacheDubboProxyService.
@@ -93,6 +105,7 @@ public final class ApacheDubboProxyServiceTest {
@AfterEach
public void after() {
ApacheDubboConfigCache.getInstance().invalidateAll();
+ RpcContext.removeClientAttachment();
}
@Test
@@ -112,6 +125,68 @@ public final class ApacheDubboProxyServiceTest {
future.complete("success");
}
+ @Test
+ @SuppressWarnings("unchecked")
+ void defersReferenceAccessAndInvocationOffTheRequestThread() throws
Exception {
+ GenericService genericService = mock(GenericService.class);
+ AtomicReference<Thread> worker = new AtomicReference<>();
+ AtomicReference<RpcContextAttachment> workerContext = new
AtomicReference<>();
+ RpcContext.getClientAttachment().setAttachment("timeout", 1234);
+
RpcContext.getClientAttachment().setAttachment(Constants.DUBBO_SELECTOR_ID,
selectorData.getId());
+
RpcContext.getClientAttachment().setAttachment(Constants.DUBBO_RULE_ID,
ruleData.getId());
+
RpcContext.getClientAttachment().setAttachment(Constants.DUBBO_REMOTE_ADDRESS,
"127.0.0.1");
+ RpcContext.getClientAttachment().setAttachment("custom", "value");
+ when(referenceConfig.getInterface()).thenReturn(PATH);
+ when(referenceConfig.get()).thenAnswer(invocation -> {
+ assertFalse(Schedulers.isInNonBlockingThread());
+ worker.set(Thread.currentThread());
+ return genericService;
+ });
+ when(genericService.$invoke(METHOD_NAME, LEFT,
RIGHT)).thenAnswer(invocation -> {
+ assertFalse(Schedulers.isInNonBlockingThread());
+ assertTrue(Thread.currentThread() == worker.get());
+ workerContext.set(RpcContext.getClientAttachment());
+ assertEquals(1234,
workerContext.get().getObjectAttachments().get("timeout"));
+ assertEquals(selectorData.getId(),
workerContext.get().getAttachment(Constants.DUBBO_SELECTOR_ID));
+ assertEquals(ruleData.getId(),
workerContext.get().getAttachment(Constants.DUBBO_RULE_ID));
+ assertEquals("127.0.0.1",
workerContext.get().getAttachment(Constants.DUBBO_REMOTE_ADDRESS));
+ assertEquals("value", workerContext.get().getAttachment("custom"));
+ return null;
+ });
+ Field field = ApacheDubboConfigCache.class.getDeclaredField("cache");
+ field.setAccessible(true);
+ ((LoadingCache<String, ReferenceConfig<GenericService>>)
field.get(ApacheDubboConfigCache.getInstance())).put(PATH, referenceConfig);
+ ApacheDubboProxyService service = new ApacheDubboProxyService(new
BodyParamResolveServiceImpl());
+ Mono<Object> result = service.genericInvoker("", metaData,
selectorData, ruleData, exchange);
+ verifyNoInteractions(referenceConfig, genericService);
+
+
StepVerifier.create(result.subscribeOn(Schedulers.parallel())).expectNext(Constants.DUBBO_RPC_RESULT_EMPTY).verifyComplete();
+
+ assertNotSame(Thread.currentThread(), worker.get());
+ assertTrue(workerContext.get().getObjectAttachments().isEmpty());
+ assertEquals("value",
RpcContext.getClientAttachment().getAttachment("custom"));
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void referenceInitializationErrorsAreReactive() throws Exception {
+ when(referenceConfig.getInterface()).thenReturn(PATH);
+ AtomicReference<RpcContextAttachment> workerContext = new
AtomicReference<>();
+ RpcContext.getClientAttachment().setAttachment("custom", "value");
+ when(referenceConfig.get()).thenAnswer(invocation -> {
+ workerContext.set(RpcContext.getClientAttachment());
+ assertEquals("value", workerContext.get().getAttachment("custom"));
+ throw new IllegalStateException("registry unavailable");
+ });
+ Field field = ApacheDubboConfigCache.class.getDeclaredField("cache");
+ field.setAccessible(true);
+ ((LoadingCache<String, ReferenceConfig<GenericService>>)
field.get(ApacheDubboConfigCache.getInstance())).put(PATH, referenceConfig);
+ ApacheDubboProxyService service = new ApacheDubboProxyService(new
BodyParamResolveServiceImpl());
+ Mono<Object> result = service.genericInvoker("", metaData,
selectorData, ruleData, exchange);
+ StepVerifier.create(result).expectErrorMessage("registry
unavailable").verify();
+ assertTrue(workerContext.get().getObjectAttachments().isEmpty());
+ }
+
static class BodyParamResolveServiceImpl implements
DubboParamResolveService {
@Override