dubbo 服务端注册流程

一、启动一个服务端Provider
1. 定义一个接口和实现
1public interface UserService {
2 void say(String message);
3}
4public class UserServiceImpl implements UserService {
5 public void say(String message) {
6 System.out.println("say:" + message);
7 }
8}2. 本地服务注册到zk
1public class Tester {
2
3 @Test
4 public void providerTest() {
5
6 //1. 服务方要把UserService方法提供给外面调用
7 UserService userService = new UserServiceImpl();
8
9 //2. 应用配置
10 ApplicationConfig app = new ApplicationConfig();
11 app.setName("providerTest");
12
13 //3. 指定一个注册中心
14 RegistryConfig registry = new RegistryConfig();
15 registry.setAddress("zookeeper://127.0.0.1:2181");
16
17 //4. 指定协议类型
18 ProtocolConfig protocol = new ProtocolConfig();
19 protocol.setName("dubbo");
20 protocol.setPort(8012);
21 protocol.setThreads(200);
22
23 // 服务提供者暴露服务配置
24 ServiceConfig<UserService> service = new ServiceConfig<UserService>(); // 此实例很重,封装了与注册中心的连接,请自行缓存,否则可能造成内存和连接泄漏
25 service.setApplication(app);
26 service.setRegistry(registry); // 多个注册中心可以用setRegistries()
27 service.setProtocol(protocol); // 多个协议可以用setProtocols()
28 service.setInterface(UserService.class);
29 service.setRef(userService);
30 service.setVersion("1.0.0");
31
32 // 暴露及注册服务
33 //dubbo://192.168.1.9:8012/code.UserService?anyhost=true&application=providerTest&dubbo=2.5.3&interface=code.UserService&methods=say&pid=46787&revision=1.0.0&side=provider&threads=200×tamp=1597048727957&version=1.0.0
34 service.export();
35
36 }
37
38
39}3. 分析原理
这里只是分析下大概原理,给各位童靴先带来带你感受,实际步骤后面分析源码时候再细说
在进行分析之前我们思考一下,当我们不使用RPC框架和SpringCloud的时候,如果我们要调用其他第三方的服务时候,我们会怎么处理呢?
通过下面这中方式每次调用时候构建一个HTTP的请求。
1public class Tester{
2 public static void sayRequest(String message){
3 OkHttpClient client = new OkHttpClient();
4 Request request = new Request.Builder()
5 .url("http://第三方服务的接口地址?message"+message)
6 .get().addHeader("Cache-Control","no-cache")
7 .build();
8 client.newCall(request).execute();
9 }
10 public static void main(String[]args){
11 sayRequest("你好")
12 }
13}
14使用后我们就可以像调用本地方法一样来调用远程接口了? 那么Dubbo是如何实现的呢? 其实就是在底层帮我们做了类似于http的通信 而通过api的方式屏蔽了底层。让我们直接将调用本地方法一样调用远程方法。
关键词一:通信协议
dubbo默认不是基于HTTP,而是基于dubbo自定义的协议。因为jdk自带的socket api不太友好,所以dubbo底层是使用netty类做通信的 说白了这个协议和http类似都是基于tcp协议从而进行封装,不同点就是数据格式不同。 如下我们自定义了一个协议来读取tcp连接中数据。
下面代码不是重点,重点知道协议就是,约定从tcp连接中读取数据的方式和方法。比如约定了读的第一个字节是协议类型,第二个是序列化类型,第三个是报文数据长度,第四个就是具体的报文数据。
1 /**
2 * 主要依据:
3 * 数据有
4 * 协议类型(1位) + 序列化类型(1位) + 报文大小(4位) + 数据报文组成(N位)
5 * *******************************************************************************
6 * ---------------- ----------------- ---------------- ------------------
7 * | 协议类型(1位) | + | 序列化类型(1位) | + | 报文大小(4位) | + | 数据报文组成(N位)|
8 * ---------------- ----------------- ---------------- ------------------
9 * *******************************************************************************¬
10 **/
11 @Override
12 public void doDecode(ChannelHandlerContext ctx, ByteBuf inByteBuf, List<Object> out) throws Exception {
13 byte[] dataArr;
14 //1. 不可读就关闭
15 if (!inByteBuf.isReadable()) {
16 Channel channel = ctx.channel();
17 SocketAddress socketAddress = channel.remoteAddress();
18 channel.close();
19 System.err.println(">>>>>>>>>[" + socketAddress + "]客户端已主动断开连接....");
20 return;
21 }
22 //2. 可读的数据大小
23 int dataHeadSize = inByteBuf.readableBytes();
24 //3. 不是完整的数据头就直接返回
25 if (!isFullMessageHeader(dataHeadSize)) {
26 return;
27 }
28 //4. 完整的数据头就开始看数据长度是否满足
29 inByteBuf.markReaderIndex();
30 //协议类型
31 byte protocolType = inByteBuf.readByte();
32 //序列化类型
33 byte serializationType = inByteBuf.readByte();
34 //数据长度
35 int dataSize = inByteBuf.readInt();
36 //5. 拆包的直接返回下次数据完整了,在处理
37 if (!isFullMessage(inByteBuf, dataSize)) {
38 inByteBuf.resetReaderIndex();
39 System.out.println();
40 System.err.println("######################数据不足已重置buffer######################");
41 return;
42 }
43 System.out.println();
44 System.err.println("######################数据完整######################");
45 //6. 黏包的直接读取数据
46 dataArr = new byte[dataSize];
47 inByteBuf.readBytes(dataArr, 0, dataSize);
48 //找到序列化器,性能有提升空间,可以序列化器可以进行池化
49 SerializeEnum serializeEnum = SerializeEnum.ofByType(serializationType);
50 Class<? extends Serialize> serialize = serializeEnum.getSerialize();
51 //根据类型获取序列化器
52 Serialize serializeNewInstance = serialize.newInstance();
53 Object deserialize = serializeNewInstance.deserialize(dataArr);
54 out.add(deserialize);
55 }关键词二:封装通信
服务端
- 服务端将需要提供的接口实现方法封装起来
- 并启动一个Netty服务
- 同时将自己的地址注册到zk中
客户端
- 客户端通过将接口方法封装成URL
- 去请求zk,拿到真实的provider地址
- 具体调用时候去请求服务端的netty服务之星

二、提供服务流程
这里我们只先分析dubbo的源码,后面再说dubbo整合spring的原理。
1. 要提供服务的对象
1public interface UserService {
2 void say(String message);
3}
4public class UserServiceImpl implements UserService {
5 public void say(String message) {
6 System.out.println("say:" + message);
7 }
8}2. 创建一个应用
1ApplicationConfig app = new ApplicationConfig();
2app.setName("providerTest");3. 指定注册中心
这里我们使用zookeeper作为注册中心
1 RegistryConfig registry = new RegistryConfig();
2 registry.setAddress("zookeeper://127.0.0.1:2181");4. 指定通信协议
1 ProtocolConfig protocol = new ProtocolConfig();
2 protocol.setName("dubbo");
3 protocol.setPort(8012);
4 protocol.setThreads(200);5. 导出服务到zk
1 // 服务提供者暴露服务配置
2 ServiceConfig<UserService> service = new ServiceConfig<UserService>(); // 此实例很重,封装了与注册中心的连接,请自行缓存,否则可能造成内存和连接泄漏
3 service.setApplication(app);
4 service.setRegistry(registry); // 多个注册中心可以用setRegistries()
5 service.setProtocol(protocol); // 多个协议可以用setProtocols()
6 service.setInterface(UserService.class);
7 service.setRef(userService);
8 service.setVersion("1.0.0");
9
10 // 暴露及注册服务
11 //dubbo://192.168.1.9:8012/code.UserService?anyhost=true&application=providerTest&dubbo=2.5.3&interface=code.UserService&methods=say&pid=46787&revision=1.0.0&side=provider&threads=200×tamp=1597048727957&version=1.0.0
12 service.export();当这一步进行完后,我们会在zookeeper的控制台找到自己的服务地址。

通过url解码之后就是
dubbo://192.168.1.9:8012/code.UserService?anyhost=true&application=providerTest&dubbo=2.5.3&interface=code.UserService&methods=say&pid=46787&revision=1.0.0&side=provider&threads=200×tamp=1597048727957&version=1.0.0
三、源码分析
1. 服务注册
我们在看二流程中,可以看到前面的1234创建的步骤都是在5中使用的,说明1234其实都是数据的载体,具体如何使用是在5中来使用的。而5的对象是ServiceConfig。所以说看源码的入口就从ServiceConfig.export()开始。
ServiceConfig的export最终后调用doExport();

doExport方法会先检查然后在注册服务到Netty服务器和注册到zk
1 protected synchronized void doExport() {
2 interfaceClass = Class.forName(interfaceName, true, Thread.currentThread()
3 .getContextClassLoader());
4 //检查接口方法是否存在在接口中,这种是使用的方法级别的执行时候
5 checkInterfaceAndMethods(interfaceClass, methods);
6 //检查接口实例是否存在,必须存在否则无法执行反射
7 checkRef();
8 //检查应该配置,如果没有配置自动创建一个,应用名是dubbo.application.name的值
9 checkApplication();
10 //检查注册中心
11 checkRegistry();
12 //检查协议,默认是dubbo协议
13 checkProtocol();
14 appendProperties(this);
15 checkStubAndMock(interfaceClass);
16 if (path == null || path.length() == 0) {
17 path = interfaceName;
18 }
19 //真正导出服务
20 doExportUrls();
21 }这一步会将服务在本地启动一个服务,同时将服务注册到注册中心中。
1 private void doExportUrls() {
2 List<URL> registryURLs = loadRegistries(true);
3 //registry://127.0.0.1:2181/com.alibaba.dubbo.registry.RegistryService?application=providerTest&dubbo=2.5.3&pid=48656®istry=zookeeper×tamp=1597052329640
4 for (ProtocolConfig protocolConfig : protocols) {
5 //核心逻辑在这里
6 doExportUrlsFor1Protocol(protocolConfig, registryURLs);
7 }
8 }根据doExportUrlsFor1Protocol的源码。我们发现dubbo中的所有模型都向Invoker来靠拢。
- 先创建一个Socket来验证下能不能连接上注册中心
- 然后根据协议信息,找到实现类。端口如果指定了就用指定的,没有指定就随机生成。默认是20880
- 获取服务版本号,首先查找MANIFEST.MF规范中的版本号。没有就用指定的版本号
- 通过反射生成Invoker对象
- 导出Invoker启动一个Netty服务DubboProtocol.openServer
- 注册到zk中RegistryProtocol.export
2. Netty服务接受服务
前面注册时候,我们说了在DubboProtocol中去创建服务的。那我们直接看这部分代码。
1public class DubboProtocol extends AbstractProtocol {
2 //逻辑处理器
3 private ExchangeHandler requestHandler = new ExchangeHandlerAdapter() {
4
5 public Object reply(ExchangeChannel channel, Object message) throws RemotingException {
6 if (message instanceof Invocation) {
7 Invocation inv = (Invocation) message;
8 Invoker<?> invoker = getInvoker(channel, inv);
9 //如果是callback 需要处理高版本调用低版本的问题
10 if (Boolean.TRUE.toString().equals(inv.getAttachments().get(IS_CALLBACK_SERVICE_INVOKE))){
11 String methodsStr = invoker.getUrl().getParameters().get("methods");
12 boolean hasMethod = false;
13 if (methodsStr == null || methodsStr.indexOf(",") == -1){
14 hasMethod = inv.getMethodName().equals(methodsStr);
15 } else {
16 String[] methods = methodsStr.split(",");
17 for (String method : methods){
18 if (inv.getMethodName().equals(method)){
19 hasMethod = true;
20 break;
21 }
22 }
23 }
24 if (!hasMethod){
25 logger.warn(new IllegalStateException("The methodName "+inv.getMethodName()+" not found in callback service interface ,invoke will be ignored. please update the api interface. url is:" + invoker.getUrl()) +" ,invocation is :"+inv );
26 return null;
27 }
28 }
29 RpcContext.getContext().setRemoteAddress(channel.getRemoteAddress());
30 return invoker.invoke(inv);
31 }
32 throw new RemotingException(channel, "Unsupported request: " + message == null ? null : (message.getClass().getName() + ": " + message) + ", channel: consumer: " + channel.getRemoteAddress() + " --> provider: " + channel.getLocalAddress());
33 }
34 //创建服务
35 private void openServer(URL url) {
36 // find server.
37 String key = url.getAddress();
38 //client 也可以暴露一个只有server可以调用的服务。
39 boolean isServer = url.getParameter(Constants.IS_SERVER_KEY,true);
40 if (isServer) {
41 ExchangeServer server = serverMap.get(key);
42 if (server == null) {
43 serverMap.put(key, createServer(url));
44 } else {
45 //server支持reset,配合override功能使用
46 server.reset(url);
47 }
48 }
49 }
50 //创建服务
51 private ExchangeServer createServer(URL url) {
52 //默认开启server关闭时发送readonly事件
53 url = url.addParameterIfAbsent(Constants.CHANNEL_READONLYEVENT_SENT_KEY, Boolean.TRUE.toString());
54 //默认开启heartbeat
55 url = url.addParameterIfAbsent(Constants.HEARTBEAT_KEY, String.valueOf(Constants.DEFAULT_HEARTBEAT));
56 String str = url.getParameter(Constants.SERVER_KEY, Constants.DEFAULT_REMOTING_SERVER);
57
58 if (str != null && str.length() > 0 && ! ExtensionLoader.getExtensionLoader(Transporter.class).hasExtension(str))
59 throw new RpcException("Unsupported server type: " + str + ", url: " + url);
60
61 url = url.addParameter(Constants.CODEC_KEY, Version.isCompatibleVersion() ? COMPATIBLE_CODEC_NAME : DubboCodec.NAME);
62 ExchangeServer server;
63 try {
64 server = Exchangers.bind(url, requestHandler);
65 } catch (RemotingException e) {
66 throw new RpcException("Fail to start server(url: " + url + ") " + e.getMessage(), e);
67 }
68 str = url.getParameter(Constants.CLIENT_KEY);
69 if (str != null && str.length() > 0) {
70 Set<String> supportedTypes = ExtensionLoader.getExtensionLoader(Transporter.class).getSupportedExtensions();
71 if (!supportedTypes.contains(str)) {
72 throw new RpcException("Unsupported client type: " + str);
73 }
74 }
75 return server;
76 }
77}我们主要看服务端要的3个方法
- ExchangeHandler逻辑处理器
- 创建服务openServer和createServer。
- 底层实现NettyTransporter
ExchangeHandler#ExchangeHandler
既然我们说了底层是Netty来实现的,那么又知道Netty是通信框架。那么我们来看下服务端的处理逻辑吧。
1class DubboProtocol{
2 private ExchangeHandler requestHandler = new ExchangeHandlerAdapter() {
3 public Object reply(ExchangeChannel channel, Object message) throws RemotingException {
4 if (message instanceof Invocation) {
5 Invocation inv = (Invocation) message;
6 Invoker<?> invoker = getInvoker(channel, inv);
7 //如果是callback 需要处理高版本调用低版本的问题
8 if (Boolean.TRUE.toString().equals(inv.getAttachments().get(IS_CALLBACK_SERVICE_INVOKE))){
9 String methodsStr = invoker.getUrl().getParameters().get("methods");
10 boolean hasMethod = false;
11 if (methodsStr == null || methodsStr.indexOf(",") == -1){
12 hasMethod = inv.getMethodName().equals(methodsStr);
13 } else {
14 String[] methods = methodsStr.split(",");
15 for (String method : methods){
16 if (inv.getMethodName().equals(method)){
17 hasMethod = true;
18 break;
19 }
20 }
21 }
22 if (!hasMethod){
23 logger.warn(new IllegalStateException("The methodName "+inv.getMethodName()+" not found in callback service interface ,invoke will be ignored. please update the api interface. url is:" + invoker.getUrl()) +" ,invocation is :"+inv );
24 return null;
25 }
26 }
27 RpcContext.getContext().setRemoteAddress(channel.getRemoteAddress());
28 return invoker.invoke(inv);
29 }
30 throw new RemotingException(channel, "Unsupported request: " + message == null ? null : (message.getClass().getName() + ": " + message) + ", channel: consumer: " + channel.getRemoteAddress() + " --> provider: " + channel.getLocalAddress());
31 }
32 }
33}- 过滤器
请看注释

3. 编码器和解码器
这里稍微说一点编码器,dubbo协议的编码器。

1public class NettyServer extends AbstractServer implements Server {
2
3
4 @Override
5 protected void doOpen() throws Throwable {
6 NettyHelper.setNettyLoggerFactory();
7 ExecutorService boss = Executors.newCachedThreadPool(new NamedThreadFactory("NettyServerBoss", true));
8 ExecutorService worker = Executors.newCachedThreadPool(new NamedThreadFactory("NettyServerWorker", true));
9 ChannelFactory channelFactory = new NioServerSocketChannelFactory(boss, worker, getUrl().getPositiveParameter(Constants.IO_THREADS_KEY, Constants.DEFAULT_IO_THREADS));
10 bootstrap = new ServerBootstrap(channelFactory);
11
12 final NettyHandler nettyHandler = new NettyHandler(getUrl(), this);
13 channels = nettyHandler.getChannels();
14 // https://issues.jboss.org/browse/NETTY-365
15 // https://issues.jboss.org/browse/NETTY-379
16 // final Timer timer = new HashedWheelTimer(new NamedThreadFactory("NettyIdleTimer", true));
17 bootstrap.setPipelineFactory(new ChannelPipelineFactory() {
18 public ChannelPipeline getPipeline() {
19 NettyCodecAdapter adapter = new NettyCodecAdapter(getCodec() ,getUrl(), NettyServer.this);
20 ChannelPipeline pipeline = Channels.pipeline();
21 /*int idleTimeout = getIdleTimeout();
22 if (idleTimeout > 10000) {
23 pipeline.addLast("timer", new IdleStateHandler(timer, idleTimeout / 1000, 0, 0));
24 }*/
25 //解码器
26 pipeline.addLast("decoder", adapter.getDecoder());
27 //编码器
28 pipeline.addLast("encoder", adapter.getEncoder());
29 pipeline.addLast("handler", nettyHandler);
30 return pipeline;
31 }
32 });
33 // bind
34 channel = bootstrap.bind(getBindAddress());
35 }
36}NettyCodecAdapter adapter = new NettyCodecAdapter(getCodec() ,getUrl(), NettyServer.this);
主要看这个类Codec2

我们主要看服务端如何将tcp二进制数据转成dubbo里面的模型。
客户端: 数据DecodeableRpcInvocation -> 通过编码器转换成 -> 二进制数据
服务端: 二进制数据 -> 解码器 -> DecodeableRpcInvocation -> DubboProtocol#requestHandler处理
4. 序列化协议
java模型如何转二进制,就是序列化协议。我们所说的hession2协议就在这里用的。这里追求的是速度快,数据小。
Serialization s = CodecSupport.getSerialization(channel.getUrl(), proto);

四、Invoker
前面说了dubbo中的模型都想Invoker靠拢。其实说白了就是反射。
1. 生成Invoker对象
可以看到dubbo里面已经提供了,构建方法。我们先熟悉如何使用其API。然后把这些小的知识点慢慢的串起来就好了。
1public class Tester {
2 @Test
3 public void buildInvokerTest() {
4 JavassistProxyFactory factory = new JavassistProxyFactory();
5 UserService userService = new UserServiceImpl();
6 URL dubboUrl = URL.valueOf("test://");
7 final Invoker<UserService> invoker = factory.getInvoker(userService, UserService.class, dubboUrl);
8 }2. 创建执行参数
Invoker是执行体,Invocation是执行参数
1public class Tester {
2 public void buildInvokerTest() {
3 final Invoker<UserService> invoker = factory.getInvoker(userService, UserService.class, dubboUrl);
4
5 Invocation invocation = new Invocation() {
6
7 public String getMethodName() {
8 return "say";
9 }
10
11 public Class<?>[] getParameterTypes() {
12 return new Class[]{String.class};
13 }
14
15 public Object[] getArguments() {
16 return new Object[]{"hello"};
17 }
18
19 public Map<String, String> getAttachments() {
20 return null;
21 }
22
23 public String getAttachment(String key) {
24 return null;
25 }
26
27 public String getAttachment(String key, String defaultValue) {
28 return null;
29 }
30
31 public Invoker<?> getInvoker() {
32 return null;
33 }
34 };
35 System.out.println(invoker.invoke(invocation));
36 }
37}3. 构建带有过滤器的Invoker
1@Test
2 public void linkedInvokerTest() {
3 JavassistProxyFactory factory = new JavassistProxyFactory();
4 UserService userService = new UserServiceImpl();
5 URL dubboUrl = URL.valueOf("test://");
6 final Invoker<UserService> invoker = factory.getInvoker(userService, UserService.class, dubboUrl);
7 List<Filter> filters = new ArrayList();
8 Filter filter = new Filter() {
9 public Result invoke(Invoker<?> invoker, Invocation invocation) throws RpcException {
10 System.out.println("-------执行过滤器-------");
11 return invoker.invoke(invocation);
12 }
13 };
14 filters.add(filter);
15 Invoker<UserService> userServiceInvoker = buildInvokerChain(invoker, filters);
16 userServiceInvoker.invoke(new Invocation() {
17 public String getMethodName() {
18 return "say";
19 }
20
21 public Class<?>[] getParameterTypes() {
22 return new Class[]{String.class};
23 }
24
25 public Object[] getArguments() {
26 return new Object[]{"hello 过滤器"};
27 }
28
29 public Map<String, String> getAttachments() {
30 return null;
31 }
32
33 public String getAttachment(String key) {
34 return null;
35 }
36
37 public String getAttachment(String key, String defaultValue) {
38 return null;
39 }
40
41 public Invoker<?> getInvoker() {
42 return null;
43 }
44 });
45 }
46五、总结
知识点回顾
- 如何将对象方法生成Invoker
- 如何将Invoker注册到注册地中心
- 如何处理客户端的请求
- 二进制数据转java数据协议
- 协议中包含的序列化知识
最后求关注,求订阅,谢谢你的阅读!
下一篇会讲,dubbo如何与Spring进行整合。