大家好,我是小趴菜,接下来我会从0到1手写一个RPC框架,该专题包括以下专题,有兴趣的小伙伴就跟着我一起学习吧
本章源码地址:gitee.com/baojh123/se…
自定义注解 -> opt-01
服务提供者收发消息基础实现 -> opt-01
自定义网络传输协议的实现 -> opt-02
自定义编解码实现 -> opt-03
服务提供者调用真实方法实现 -> opt-04
完善服务消费者发送消息基础功能 -> opt-05
注册中心基础功能实现 -> opt-06
服务提供者整合注册中心 -> opt-07
服务消费者整合注册中心 -> opt-08
完善服务消费者接收响应结果 -> opt-09
服务消费者,服务提供者整合SpringBoot -> opt-10
动态代理屏蔽RPC服务调用底层细节 -> opt-10
SPI机制基础功能实现 -> opt-11
SPI机制扩展随机负载均衡策略 -> opt-12
SPI机制扩展轮询负载均衡策略 -> opt-13
SPI机制扩展JDK序列化 -> opt-14
SPI机制扩展JSON序列化 -> opt-15
SPI机制扩展protustuff序列化 -> opt-16
前言
在之前我们已经完成服务提供者,服务消费者,注册中心的相关功能,以及完成了服务提供者的收发消息。
但是现在回看服务消费者
@GetMapping("/consumer")
public String test() throws Exception{
RpcConsumer rpcConsumer = new RpcConsumer();
//这里没有获取到返回结果
rpcConsumer.sendRequest(getRequest());
Thread.sleep(5000);
rpcConsumer.close();
return "ok";
}
所以这一章我们要完成服务消费者的接收响应结果
实现
修改 com.xpc.rpc.consumer.handler.RpcConsumerHandler,也就是服务消费者的自定义处理器
添加一个成员变量
//缓存请求对应的响应结果
private Map handleResultMap = new ConcurrentHashMap();
修改 channelRead0() 方法
@Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, ProtocolMessage protocolMessage) throws Exception {
// 从服务端接收响应数据,并保存到缓存结果当中
handleResultMap.put(protocolMessage.getRpcHeader().getRequestId(),protocolMessage.getT().getData());
}
修改发送请求的方法:sendRequest()
public Object sendRequest(ProtocolMessage requestProtocolMessage) {
channel.writeAndFlush(requestProtocolMessage);
//这里死循环来获取响应结果
while(true) {
Object result = handleResultMap.remove(requestProtocolMessage.getRpcHeader().getRequestId());
if(result != null) {
return result;
}
}
}
修改 com.xpc.rpc.consumer.RpcConsumer的sendRequest()方法
public Object sendRequest(ProtocolMessage requestProtocolMessage) {
ServiceMeta serviceMeta = null;
try {
serviceMeta = registerService.discovery(requestProtocolMessage.getT().getClassName());
} catch (Exception e) {
e.printStackTrace();
}
RpcRequest request = requestProtocolMessage.getT();
String key = buildHandlerMapKey(request.getClassName(), request.getMethodName());
RpcConsumerHandler consumerHandler;
if(handlerMap.containsKey(key)) {
consumerHandler = handlerMap.get(key);
Channel channel = consumerHandler.getChannel();
if(!channel.isOpen() || !channel.isActive()) {
consumerHandler = getConsumerHandler(key, serviceMeta.getServerAddress(), serviceMeta.getRegisterPort());
handlerMap.put(buildHandlerMapKey(request.getClassName(),request.getMethodName()),consumerHandler);
}
}else {
consumerHandler = getConsumerHandler(key, serviceMeta.getServerAddress(), serviceMeta.getRegisterPort());
if(consumerHandler == null) {
throw new RuntimeException("");
}
handlerMap.put(buildHandlerMapKey(request.getClassName(),request.getMethodName()),consumerHandler);
}
//这里返回响应结果
return consumerHandler.sendRequest(requestProtocolMessage);
}
测试
修改服务消费者服务 xpc-rpc-web-consumer
package com.xpc.consumer;
import com.xpc.rpc.common.enums.RpcMsgType;
import com.xpc.rpc.consumer.RpcConsumer;
import com.xpc.rpc.protocol.ProtocolMessage;
import com.xpc.rpc.protocol.header.RpcHeader;
import com.xpc.rpc.protocol.request.RpcRequest;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
@SpringBootApplication
public class App {
private static final Logger LOGGER = LoggerFactory.getLogger(App.class);
public static void main(String[] args) {
SpringApplication.run(App.class,args);
}
@GetMapping("/consumer")
public String test() throws Exception{
RpcConsumer rpcConsumer = new RpcConsumer();
//获取响应结果
Object result = rpcConsumer.sendRequest(getRequest());
LOGGER.info("返回结果是:{}",result);
rpcConsumer.close();
return "ok";
}
private ProtocolMessage getRequest() {
ProtocolMessage protocolMessage = new ProtocolMessage();
RpcHeader rpcHeader = new RpcHeader();
rpcHeader.setMsgType(RpcMsgType.REQUEST.getType());
rpcHeader.setRequestId(1L);
RpcRequest rpcRequest = new RpcRequest();
rpcRequest.setClassName("com.xpc.interfaces.UserService");
rpcRequest.setMethodName("hello");
rpcRequest.setParameterTypes(new Class[]{String.class});
rpcRequest.setParameters(new Object[]{"coco"});
protocolMessage.setRpcHeader(rpcHeader);
protocolMessage.setT(rpcRequest);
return protocolMessage;
}
}
先启动服务提供者服务 xpc-rpc-web-provider,然后在浏览器输入:localhost:8080/test 来启动Netty服务
然后启动 xpc-rpc-web-consumer服务,然后在浏览器输入:localhost:8080/consumer
可以正常返回响应结果,服务消费者也能正常获取