dubbo实现同时请求所有节点并返回
实现该功能需要先了解dubbo的运行流程,在请求的时候需要获取所有可用的服务invoker,通过实现LoadBalance和AbstractClusterInvoker都能获取到所有的invoker,但是LoadBalance是返回一个合适的invoker,但是AbstractClusterInvoker是容错机制的实现父类,此处会针对LoadBalance返回的invoker执行调用,如果失败在去执行对应的策略,所以我们可以对AbstractClusterInvoker方法进行重新。
在源码中有一个ForkingCluster类 也是会全部调用一遍,但是只是返回最快的一个请求结果,所以我们实现一个自己的类
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.Directory;
import org.apache.dubbo.rpc.cluster.support.AbstractClusterInvoker;
import org.apache.dubbo.rpc.cluster.support.wrapper.AbstractCluster;
public class MyForking extends AbstractCluster {
//命名为myForking
public final static String NAME = "myForking";
@Override
public <T> AbstractClusterInvoker<T> doJoin(Directory<T> directory) throws RpcException {
return new MyForkingClusterInvoker<T>(directory);
}
}
import com.google.common.collect.Lists;
import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.Result;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.Directory;
import org.apache.dubbo.rpc.cluster.LoadBalance;
import org.apache.dubbo.rpc.cluster.support.AbstractClusterInvoker;
import java.util.List;
public class MyForkingClusterInvoker<T> extends AbstractClusterInvoker<T> {
public MyForkingClusterInvoker(Directory<T> directory) {
super(directory);
}
@Override
protected Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
//存放所有集群节点的返回结果
List list = Lists.newArrayList();
Result result = null;
//循环所有集群节点的invoker
for (Invoker<T> invoker : invokers) {
//执行调用
Result innerResult = invokeWithContext(invoker, invocation);
if (null == result) {
//初始化Result 用于返回
result = innerResult;
}
//将每个返回结果添加到集合中
list.add(result.getValue());
}
//修改Result的结果为list
result.setValue(list);
return result;
}
}
配置我们的自定义容错类
创建 org.apache.dubbo.rpc.cluster.Cluster 文件并添加内容:
myForking=com.jiangzheng.course.dubbo.consumer.loadbalance.MyForkingCluster
我们可以测试下
import com.alibaba.dubbo.common.Constants;
import org.apache.dubbo.common.extension.ExtensionLoader;
import org.apache.dubbo.common.io.UnsafeByteArrayInputStream;
import org.apache.dubbo.common.io.UnsafeByteArrayOutputStream;
import org.apache.dubbo.common.serialize.Serialization;
import org.apache.dubbo.config.ApplicationConfig;
import org.apache.dubbo.config.ReferenceConfig;
import org.apache.dubbo.config.RegistryConfig;
import org.apache.dubbo.rpc.service.GenericService;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.ComponentScan;
import java.io.IOException;
import java.util.List;
@SpringBootApplication
@ComponentScan(value = {"com.jiangzheng.course.dubbo.consumer.controller"})
public class Application {
public static void main(String[] args) throws IOException, ClassNotFoundException {
//启动springboot
SpringApplication.run(DubboDemoApplication.class, args);
//测试
test();
}
public static void test() throws ClassNotFoundException, IOException {
// 1.泛型参数固定为GenericService
ReferenceConfig<GenericService> referenceConfig = new ReferenceConfig<GenericService>();
referenceConfig.setApplication(new ApplicationConfig("dubbo-consumer"));
referenceConfig.setRegistry(new RegistryConfig("zookeeper://127.0.0.1:2181"));
//配置我们的容错机制为 我们自定义的 myForking
referenceConfig.setCluster("myForking");
// 2. 设置为泛化引用,并且泛化类型为nativejava
referenceConfig.setInterface("com.jiangzheng.course.dubbo.api.service.ServiceDemo");
referenceConfig.setGeneric("nativejava");
GenericService greetingService = referenceConfig.get();
UnsafeByteArrayOutputStream out = new UnsafeByteArrayOutputStream();
// 4.泛型调用, 需要把参数使用Java序列化为二进制
ExtensionLoader.getExtensionLoader(Serialization.class)
.getExtension(Constants.GENERIC_SERIALIZATION_NATIVE_JAVA).serialize(null, out).writeObject("world");
Object result = greetingService.$invoke("getSelf", new String[]{"java.lang.String"},
new Object[]{out.toByteArray()});
// 5.打印结果,需要把二进制结果使用Java反序列为对象
for (Object o : (List) result) {
UnsafeByteArrayInputStream in = new UnsafeByteArrayInputStream((byte[]) o);
System.out.println(ExtensionLoader.getExtensionLoader(Serialization.class)
.getExtension(Constants.GENERIC_SERIALIZATION_NATIVE_JAVA).deserialize(null, in).readObject());
}
}
}
返回结果:
world
world