Introduce Invocation#getInvokedInvokers to get invokers for ClusterFilter (#11359)
This commit is contained in:
parent
d2eb0ef458
commit
cbcbac9ce6
|
|
@ -16,6 +16,11 @@
|
|||
*/
|
||||
package org.apache.dubbo.rpc.cluster.support;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.ThreadLocalRandom;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import org.apache.dubbo.common.URL;
|
||||
import org.apache.dubbo.common.Version;
|
||||
import org.apache.dubbo.common.config.Configuration;
|
||||
|
|
@ -39,11 +44,6 @@ import org.apache.dubbo.rpc.model.ApplicationModel;
|
|||
import org.apache.dubbo.rpc.model.ScopeModelUtil;
|
||||
import org.apache.dubbo.rpc.support.RpcUtils;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.ThreadLocalRandom;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.DEFAULT_LOADBALANCE;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.DEFAULT_RESELECT_COUNT;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.ENABLE_CONNECTIVITY_VALIDATION;
|
||||
|
|
@ -376,6 +376,7 @@ public abstract class AbstractClusterInvoker<T> implements ClusterInvoker<T> {
|
|||
if (ProfilerSwitch.isEnableSimpleProfiler()) {
|
||||
InvocationProfilerUtils.enterProfiler(invocation, "Invoker invoke. Target Address: " + invoker.getUrl().getAddress());
|
||||
}
|
||||
invocation.addInvokedInvoker(invoker);
|
||||
result = invoker.invoke(invocation);
|
||||
} finally {
|
||||
clearContext(invoker);
|
||||
|
|
|
|||
|
|
@ -16,15 +16,16 @@
|
|||
*/
|
||||
package org.apache.dubbo.rpc.cluster.directory;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.apache.dubbo.rpc.AttachmentsAdapter;
|
||||
import org.apache.dubbo.rpc.Invocation;
|
||||
import org.apache.dubbo.rpc.Invoker;
|
||||
import org.apache.dubbo.rpc.model.ServiceModel;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.DUBBO_VERSION_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.GROUP_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.PATH_KEY;
|
||||
|
|
@ -169,6 +170,16 @@ class MockDirInvocation implements Invocation {
|
|||
return (String) getObjectAttachment(key, defaultValue);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addInvokedInvoker(Invoker<?> invoker) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<Invoker<?>> getInvokedInvokers() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object getObjectAttachment(String key, Object defaultValue) {
|
||||
Object result = attachments.get(key);
|
||||
|
|
|
|||
|
|
@ -16,6 +16,14 @@
|
|||
*/
|
||||
package org.apache.dubbo.rpc.cluster.support;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import org.apache.dubbo.common.URL;
|
||||
import org.apache.dubbo.common.extension.ExtensionLoader;
|
||||
import org.apache.dubbo.common.utils.NetUtils;
|
||||
|
|
@ -33,7 +41,6 @@ import org.apache.dubbo.rpc.cluster.filter.DemoService;
|
|||
import org.apache.dubbo.rpc.cluster.loadbalance.LeastActiveLoadBalance;
|
||||
import org.apache.dubbo.rpc.cluster.loadbalance.RandomLoadBalance;
|
||||
import org.apache.dubbo.rpc.cluster.loadbalance.RoundRobinLoadBalance;
|
||||
|
||||
import org.junit.jupiter.api.AfterAll;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
|
|
@ -43,13 +50,6 @@ import org.junit.jupiter.api.Disabled;
|
|||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.Mockito;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.CONSUMER;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.ENABLE_CONNECTIVITY_VALIDATION;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.INTERFACE_KEY;
|
||||
|
|
@ -206,6 +206,24 @@ class AbstractClusterInvokerTest {
|
|||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void testSelectedInvokers() {
|
||||
cluster = new AbstractClusterInvoker(dic) {
|
||||
@Override
|
||||
protected Result doInvoke(Invocation invocation, List invokers, LoadBalance loadbalance)
|
||||
throws RpcException {
|
||||
checkInvokers(invokers, invocation);
|
||||
Invoker invoker = select(loadbalance, invocation, invokers, null);
|
||||
return invokeWithContext(invoker, invocation);
|
||||
}
|
||||
};
|
||||
|
||||
// invoke
|
||||
cluster.invoke(invocation);
|
||||
|
||||
Assertions.assertEquals(Collections.singletonList(invoker1), invocation.getInvokedInvokers());
|
||||
}
|
||||
|
||||
@Test
|
||||
void testSelect_Invokersize1() throws Exception {
|
||||
invokers.clear();
|
||||
|
|
|
|||
|
|
@ -17,13 +17,14 @@
|
|||
|
||||
package com.alibaba.dubbo.rpc;
|
||||
|
||||
import org.apache.dubbo.rpc.model.ServiceModel;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.apache.dubbo.rpc.model.ServiceModel;
|
||||
|
||||
@Deprecated
|
||||
public interface Invocation extends org.apache.dubbo.rpc.Invocation {
|
||||
|
||||
|
|
@ -216,5 +217,15 @@ public interface Invocation extends org.apache.dubbo.rpc.Invocation {
|
|||
public org.apache.dubbo.rpc.Invocation getOriginal() {
|
||||
return delegate;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addInvokedInvoker(org.apache.dubbo.rpc.Invoker<?> invoker) {
|
||||
delegate.addInvokedInvoker(invoker);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<org.apache.dubbo.rpc.Invoker<?>> getInvokedInvokers() {
|
||||
return delegate.getInvokedInvokers();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -17,15 +17,16 @@
|
|||
|
||||
package com.alibaba.dubbo.rpc;
|
||||
|
||||
import com.alibaba.dubbo.common.Constants;
|
||||
import com.alibaba.dubbo.common.URL;
|
||||
|
||||
import java.io.Serializable;
|
||||
import java.lang.reflect.Method;
|
||||
import java.util.Arrays;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import com.alibaba.dubbo.common.Constants;
|
||||
import com.alibaba.dubbo.common.URL;
|
||||
|
||||
public class RpcInvocation implements Invocation, Serializable {
|
||||
|
||||
private static final long serialVersionUID = -4355285085441097045L;
|
||||
|
|
@ -198,6 +199,16 @@ public class RpcInvocation implements Invocation, Serializable {
|
|||
return value;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addInvokedInvoker(org.apache.dubbo.rpc.Invoker<?> invoker) {
|
||||
throw new UnsupportedOperationException();
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<org.apache.dubbo.rpc.Invoker<?>> getInvokedInvokers() {
|
||||
throw new UnsupportedOperationException();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "RpcInvocation [methodName=" + methodName + ", parameterTypes="
|
||||
|
|
|
|||
|
|
@ -17,17 +17,18 @@
|
|||
|
||||
package org.apache.dubbo.cache;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import org.apache.dubbo.rpc.RpcInvocation;
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import com.alibaba.dubbo.cache.Cache;
|
||||
import com.alibaba.dubbo.cache.CacheFactory;
|
||||
import com.alibaba.dubbo.common.URL;
|
||||
import com.alibaba.dubbo.rpc.Invocation;
|
||||
import com.alibaba.dubbo.rpc.Invoker;
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Map;
|
||||
|
||||
class CacheTest {
|
||||
|
||||
|
|
@ -107,5 +108,15 @@ class CacheTest {
|
|||
public Map<Object, Object> getAttributes() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addInvokedInvoker(org.apache.dubbo.rpc.Invoker<?> invoker) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<org.apache.dubbo.rpc.Invoker<?>> getInvokedInvokers() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,12 +16,13 @@
|
|||
*/
|
||||
package org.apache.dubbo.filter;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import com.alibaba.dubbo.rpc.Invocation;
|
||||
import com.alibaba.dubbo.rpc.Invoker;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.DUBBO_VERSION_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.GROUP_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.PATH_KEY;
|
||||
|
|
@ -100,4 +101,13 @@ public class LegacyInvocation implements Invocation {
|
|||
return getAttachments().get(key);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addInvokedInvoker(org.apache.dubbo.rpc.Invoker<?> invoker) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<org.apache.dubbo.rpc.Invoker<?>> getInvokedInvokers() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,15 +16,16 @@
|
|||
*/
|
||||
package org.apache.dubbo.service;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.apache.dubbo.rpc.AttachmentsAdapter;
|
||||
import org.apache.dubbo.rpc.Invocation;
|
||||
import org.apache.dubbo.rpc.Invoker;
|
||||
import org.apache.dubbo.rpc.model.ServiceModel;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.DUBBO_VERSION_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.GROUP_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.PATH_KEY;
|
||||
|
|
@ -179,4 +180,14 @@ public class MockInvocation implements Invocation {
|
|||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addInvokedInvoker(Invoker<?> invoker) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<Invoker<?>> getInvokedInvokers() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,15 +16,16 @@
|
|||
*/
|
||||
package org.apache.dubbo.rpc;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import org.apache.dubbo.common.Experimental;
|
||||
import org.apache.dubbo.rpc.model.ModuleModel;
|
||||
import org.apache.dubbo.rpc.model.ScopeModelUtil;
|
||||
import org.apache.dubbo.rpc.model.ServiceModel;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
/**
|
||||
* Invocation. (API, Prototype, NonThreadSafe)
|
||||
*
|
||||
|
|
@ -162,4 +163,20 @@ public interface Invocation {
|
|||
Object get(Object key);
|
||||
|
||||
Map<Object, Object> getAttributes();
|
||||
|
||||
/**
|
||||
* To add invoked invokers into invocation. Can be used in ClusterFilter or Filter for tracing or debugging purpose.
|
||||
* Currently, only support in consumer side.
|
||||
*
|
||||
* @param invoker invoked invokers
|
||||
*/
|
||||
void addInvokedInvoker(Invoker<?> invoker);
|
||||
|
||||
/**
|
||||
* Get all invoked invokers in current invocation.
|
||||
* NOTICE: A curtain invoker could be invoked for twice or more if retries.
|
||||
*
|
||||
* @return invokers
|
||||
*/
|
||||
List<Invoker<?>> getInvokedInvokers();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,6 +16,22 @@
|
|||
*/
|
||||
package org.apache.dubbo.rpc;
|
||||
|
||||
import java.io.Serializable;
|
||||
import java.lang.reflect.Method;
|
||||
import java.lang.reflect.Type;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.concurrent.locks.Lock;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import org.apache.dubbo.common.URL;
|
||||
import org.apache.dubbo.common.utils.ReflectUtils;
|
||||
import org.apache.dubbo.common.utils.StringUtils;
|
||||
|
|
@ -26,20 +42,6 @@ import org.apache.dubbo.rpc.model.ServiceDescriptor;
|
|||
import org.apache.dubbo.rpc.model.ServiceModel;
|
||||
import org.apache.dubbo.rpc.support.RpcUtils;
|
||||
|
||||
import java.io.Serializable;
|
||||
import java.lang.reflect.Method;
|
||||
import java.lang.reflect.Type;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.concurrent.locks.Lock;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.APPLICATION_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.GROUP_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.INTERFACE_KEY;
|
||||
|
|
@ -92,6 +94,8 @@ public class RpcInvocation implements Invocation, Serializable {
|
|||
|
||||
private transient InvokeMode invokeMode;
|
||||
|
||||
private transient List<Invoker<?>> invokedInvokers = new LinkedList<>();
|
||||
|
||||
/**
|
||||
* @deprecated only for test
|
||||
*/
|
||||
|
|
@ -311,6 +315,16 @@ public class RpcInvocation implements Invocation, Serializable {
|
|||
return attributes;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addInvokedInvoker(Invoker<?> invoker) {
|
||||
this.invokedInvokers.add(invoker);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<Invoker<?>> getInvokedInvokers() {
|
||||
return this.invokedInvokers;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getTargetServiceUniqueName() {
|
||||
return targetServiceUniqueName;
|
||||
|
|
|
|||
|
|
@ -16,10 +16,12 @@
|
|||
*/
|
||||
package org.apache.dubbo.rpc;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.HashMap;
|
||||
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.HashMap;
|
||||
import org.mockito.Mockito;
|
||||
|
||||
class RpcInvocationTest {
|
||||
|
||||
|
|
@ -43,4 +45,20 @@ class RpcInvocationTest {
|
|||
invocation.setObjectAttachments(map);
|
||||
Assertions.assertEquals(map, invocation.getObjectAttachments());
|
||||
}
|
||||
|
||||
@Test
|
||||
void testInvokers() {
|
||||
RpcInvocation rpcInvocation = new RpcInvocation();
|
||||
|
||||
Invoker invoker1 = Mockito.mock(Invoker.class);
|
||||
Invoker invoker2 = Mockito.mock(Invoker.class);
|
||||
Invoker invoker3 = Mockito.mock(Invoker.class);
|
||||
|
||||
rpcInvocation.addInvokedInvoker(invoker1);
|
||||
rpcInvocation.addInvokedInvoker(invoker2);
|
||||
rpcInvocation.addInvokedInvoker(invoker3);
|
||||
rpcInvocation.addInvokedInvoker(invoker3);
|
||||
|
||||
Assertions.assertEquals(Arrays.asList(invoker1, invoker2, invoker3, invoker3), rpcInvocation.getInvokedInvokers());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue