dubbo

dubbo 服务端注册流程

2026-01-298 min read手写RPC框架
Dubbo

一、启动一个服务端Provider

1. 定义一个接口和实现

java
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

java
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&timestamp=1597048727957&version=1.0.0 34 service.export(); 35 36 } 37 38 39}

3. 分析原理

这里只是分析下大概原理,给各位童靴先带来带你感受,实际步骤后面分析源码时候再细说

在进行分析之前我们思考一下,当我们不使用RPC框架和SpringCloud的时候,如果我们要调用其他第三方的服务时候,我们会怎么处理呢?

通过下面这中方式每次调用时候构建一个HTTP的请求。

java
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连接中读取数据的方式和方法。比如约定了读的第一个字节是协议类型,第二个是序列化类型,第三个是报文数据长度,第四个就是具体的报文数据。

java
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 }

关键词二:封装通信

服务端
  1. 服务端将需要提供的接口实现方法封装起来
  2. 并启动一个Netty服务
  3. 同时将自己的地址注册到zk中
客户端
  1. 客户端通过将接口方法封装成URL
  2. 去请求zk,拿到真实的provider地址
  3. 具体调用时候去请求服务端的netty服务之星

二、提供服务流程

这里我们只先分析dubbo的源码,后面再说dubbo整合spring的原理。

1. 要提供服务的对象

java
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. 创建一个应用

java
1ApplicationConfig app = new ApplicationConfig(); 2app.setName("providerTest");

3. 指定注册中心

这里我们使用zookeeper作为注册中心

java
1 RegistryConfig registry = new RegistryConfig(); 2 registry.setAddress("zookeeper://127.0.0.1:2181");

4. 指定通信协议

java
1 ProtocolConfig protocol = new ProtocolConfig(); 2 protocol.setName("dubbo"); 3 protocol.setPort(8012); 4 protocol.setThreads(200);

5. 导出服务到zk

java
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&timestamp=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&timestamp=1597048727957&version=1.0.0

三、源码分析

1. 服务注册

我们在看二流程中,可以看到前面的1234创建的步骤都是在5中使用的,说明1234其实都是数据的载体,具体如何使用是在5中来使用的。而5的对象是ServiceConfig。所以说看源码的入口就从ServiceConfig.export()开始。

ServiceConfig的export最终后调用doExport();

doExport方法会先检查然后在注册服务到Netty服务器和注册到zk

java
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 }

这一步会将服务在本地启动一个服务,同时将服务注册到注册中心中。

java
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&registry=zookeeper&timestamp=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中去创建服务的。那我们直接看这部分代码。

java
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是通信框架。那么我们来看下服务端的处理逻辑吧。

java
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协议的编码器。

java
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。然后把这些小的知识点慢慢的串起来就好了。

java
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是执行参数

java
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

java
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进行整合。