前段时间看到一篇不错的文章《看了这篇你就会手写RPC框架了》,于是便来了兴趣对着实现了一遍,后面觉得还有很多优化的地方便对其进行了改进。
主要改动点如下:
RPC,即 Remote Procedure Call(远程过程调用),调用远程计算机上的服务,就像调用本地服务一样。RPC可以很好的解耦系统,如WebService就是一种基于Http协议的RPC。
总的来说,就如下几个步骤:
所以一个RPC框架有如下角色:
本RPC框架rpc-spring-boot-starter涉及技术栈如下:
由于代码过多,这里只讲几处改动点。
1.编写LoadBalance的实现类
2.自定义注解 @LoadBalanceAno
/**
* 负载均衡注解
*/
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Documented
public @interface LoadBalanceAno {
String value() default "";
}
/**
* 轮询算法
*/
@LoadBalanceAno(RpcConstant.BALANCE_ROUND)
public class FullRoundBalance implements LoadBalance {
private static Logger logger = LoggerFactory.getLogger(FullRoundBalance.class);
private volatile int index;
@Override
public synchronized Service chooseOne(List<Service> services) {
// 加锁防止多线程情况下,index超出services.size()
if (index == services.size()) {
index = 0;
}
return services.get(index++);
}
}
3.新建在resource目录下META-INF/servers文件夹并创建文件 4.RpcConfig增加配置项loadBalance
/**
* @author 2YSP
* @date 2020/7/26 15:13
*/
@ConfigurationProperties(prefix = "sp.rpc")
public class RpcConfig {
/**
* 服务注册中心地址
*/
private String registerAddress = "127.0.0.1:2181";
/**
* 服务暴露端口
*/
private Integer serverPort = 9999;
/**
* 服务协议
*/
private String protocol = "java";
/**
* 负载均衡算法
*/
private String loadBalance = "random";
/**
* 权重,默认为1
*/
private Integer weight = 1;
// 省略getter setter
}
5.在自动配置类RpcAutoConfiguration根据配置选择对应的算法实现类
/**
* 使用spi匹配符合配置的负载均衡算法
*
* @param name
* @return
*/
private LoadBalance getLoadBalance(String name) {
ServiceLoader<LoadBalance> loader = ServiceLoader.load(LoadBalance.class);
Iterator<LoadBalance> iterator = loader.iterator();
while (iterator.hasNext()) {
LoadBalance loadBalance = iterator.next();
LoadBalanceAno ano = loadBalance.getClass().getAnnotation(LoadBalanceAno.class);
Assert.notNull(ano, "load balance name can not be empty!");
if (name.equals(ano.value())) {
return loadBalance;
}
}
throw new RpcException("invalid load balance config");
}
@Bean
public ClientProxyFactory proxyFactory(@Autowired RpcConfig rpcConfig) {
ClientProxyFactory clientProxyFactory = new ClientProxyFactory();
// 设置服务发现着
clientProxyFactory.setServerDiscovery(new ZookeeperServerDiscovery(rpcConfig.getRegisterAddress()));
// 设置支持的协议
Map<String, MessageProtocol> supportMessageProtocols = buildSupportMessageProtocols();
clientProxyFactory.setSupportMessageProtocols(supportMessageProtocols);
// 设置负载均衡算法
LoadBalance loadBalance = getLoadBalance(rpcConfig.getLoadBalance());
clientProxyFactory.setLoadBalance(loadBalance);
// 设置网络层实现
clientProxyFactory.setNetClient(new NettyNetClient());
return clientProxyFactory;
}
使用Map来缓存数据
/**
* 服务发现本地缓存
*/
public class ServerDiscoveryCache {
/**
* key: serviceName
*/
private static final Map<String, List<Service>> SERVER_MAP = new ConcurrentHashMap<>();
/**
* 客户端注入的远程服务service class
*/
public static final List<String> SERVICE_CLASS_NAMES = new ArrayList<>();
public static void put(String serviceName, List<Service> serviceList) {
SERVER_MAP.put(serviceName, serviceList);
}
/**
* 去除指定的值
* @param serviceName
* @param service
*/
public static void remove(String serviceName, Service service) {
SERVER_MAP.computeIfPresent(serviceName, (key, value) ->
value.stream().filter(o -> !o.toString().equals(service.toString())).collect(Collectors.toList())
);
}
public static void removeAll(String serviceName) {
SERVER_MAP.remove(serviceName);
}
public static boolean isEmpty(String serviceName) {
return SERVER_MAP.get(serviceName) == null || SERVER_MAP.get(serviceName).size() == 0;
}
public static List<Service> get(String serviceName) {
return SERVER_MAP.get(serviceName);
}
}
ClientProxyFactory,先查本地缓存,缓存没有再查询zookeeper。
/**
* 根据服务名获取可用的服务地址列表
* @param serviceName
* @return
*/
private List<Service> getServiceList(String serviceName) {
List<Service> services;
synchronized (serviceName){
if (ServerDiscoveryCache.isEmpty(serviceName)) {
services = serverDiscovery.findServiceList(serviceName);
if (services == null || services.size() == 0) {
throw new RpcException("No provider available!");
}
ServerDiscoveryCache.put(serviceName, services);
} else {
services = ServerDiscoveryCache.get(serviceName);
}
}
return services;
}
问题:如果服务端因为宕机或网络问题下线了,缓存却还在就会导致客户端请求已经不可用的服务端,增加请求失败率。解决方案:由于服务端注册的是临时节点,所以如果服务端下线节点会被移除。只要监听zookeeper的子节点,如果新增或删除子节点就直接清空本地缓存即可。
DefaultRpcProcessor
/**
* Rpc处理者,支持服务启动暴露,自动注入Service
* @author 2YSP
* @date 2020/7/26 14:46
*/
public class DefaultRpcProcessor implements ApplicationListener<ContextRefreshedEvent> {
@Override
public void onApplicationEvent(ContextRefreshedEvent event) {
// Spring启动完毕过后会收到一个事件通知
if (Objects.isNull(event.getApplicationContext().getParent())){
ApplicationContext context = event.getApplicationContext();
// 开启服务
startServer(context);
// 注入Service
injectService(context);
}
}
private void injectService(ApplicationContext context) {
String[] names = context.getBeanDefinitionNames();
for(String name : names){
Class<?> clazz = context.getType(name);
if (Objects.isNull(clazz)){
continue;
}
Field[] declaredFields = clazz.getDeclaredFields();
for(Field field : declaredFields){
// 找出标记了InjectService注解的属性
InjectService injectService = field.getAnnotation(InjectService.class);
if (injectService == null){
continue;
}
Class<?> fieldClass = field.getType();
Object object = context.getBean(name);
field.setAccessible(true);
try {
field.set(object,clientProxyFactory.getProxy(fieldClass));
} catch (IllegalAccessException e) {
e.printStackTrace();
}
// 添加本地服务缓存
ServerDiscoveryCache.SERVICE_CLASS_NAMES.add(fieldClass.getName());
}
}
// 注册子节点监听
if (clientProxyFactory.getServerDiscovery() instanceof ZookeeperServerDiscovery){
ZookeeperServerDiscovery serverDiscovery = (ZookeeperServerDiscovery) clientProxyFactory.getServerDiscovery();
ZkClient zkClient = serverDiscovery.getZkClient();
ServerDiscoveryCache.SERVICE_CLASS_NAMES.forEach(name ->{
String servicePath = RpcConstant.ZK_SERVICE_PATH + RpcConstant.PATH_DELIMITER + name + "/service";
zkClient.subscribeChildChanges(servicePath, new ZkChildListenerImpl());
});
logger.info("subscribe service zk node successfully");
}
}
private void startServer(ApplicationContext context) {
...
}
}
ZkChildListenerImpl
/**
* 子节点事件监听处理类
*/
public class ZkChildListenerImpl implements IZkChildListener {
private static Logger logger = LoggerFactory.getLogger(ZkChildListenerImpl.class);
/**
* 监听子节点的删除和新增事件
* @param parentPath /rpc/serviceName/service
* @param childList
* @throws Exception
*/
@Override
public void handleChildChange(String parentPath, List<String> childList) throws Exception {
logger.debug("Child change parentPath:[{}] -- childList:[{}]", parentPath, childList);
// 只要子节点有改动就清空缓存
String[] arr = parentPath.split("/");
ServerDiscoveryCache.removeAll(arr[2]);
}
}
这部分的改动最多,先增加新的sendRequest接口。
添加接口实现类NettyNetClient
/**
* @author 2YSP
* @date 2020/7/25 20:12
*/
public class NettyNetClient implements NetClient {
private static Logger logger = LoggerFactory.getLogger(NettyNetClient.class);
private static ExecutorService threadPool = new ThreadPoolExecutor(4, 10, 200,
TimeUnit.SECONDS, new LinkedBlockingQueue<>(1000), new ThreadFactoryBuilder()
.setNameFormat("rpcClient-%d")
.build());
private EventLoopGroup loopGroup = new NioEventLoopGroup(4);
/**
* 已连接的服务缓存
* key: 服务地址,格式:ip:port
*/
public static Map<String, SendHandlerV2> connectedServerNodes = new ConcurrentHashMap<>();
@Override
public byte[] sendRequest(byte[] data, Service service) throws InterruptedException {
....
return respData;
}
@Override
public RpcResponse sendRequest(RpcRequest rpcRequest, Service service, MessageProtocol messageProtocol) {
String address = service.getAddress();
synchronized (address) {
if (connectedServerNodes.containsKey(address)) {
SendHandlerV2 handler = connectedServerNodes.get(address);
logger.info("使用现有的连接");
return handler.sendRequest(rpcRequest);
}
String[] addrInfo = address.split(":");
final String serverAddress = addrInfo[0];
final String serverPort = addrInfo[1];
final SendHandlerV2 handler = new SendHandlerV2(messageProtocol, address);
threadPool.submit(() -> {
// 配置客户端
Bootstrap b = new Bootstrap();
b.group(loopGroup).channel(NioSocketChannel.class)
.option(ChannelOption.TCP_NODELAY, true)
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel socketChannel) throws Exception {
ChannelPipeline pipeline = socketChannel.pipeline();
pipeline
.addLast(handler);
}
});
// 启用客户端连接
ChannelFuture channelFuture = b.connect(serverAddress, Integer.parseInt(serverPort));
channelFuture.addListener(new ChannelFutureListener() {
@Override
public void operationComplete(ChannelFuture channelFuture) throws Exception {
connectedServerNodes.put(address, handler);
}
});
}
);
logger.info("使用新的连接。。。");
return handler.sendRequest(rpcRequest);
}
}
}
每次请求都会调用sendRequest()方法,用线程池异步和服务端创建TCP长连接,连接成功后将SendHandlerV2缓存到ConcurrentHashMap中方便复用,后续请求的请求地址(ip+port)如果在connectedServerNodes中存在则使用connectedServerNodes中的handler处理不再重新建立连接。
SendHandlerV2
/**
* @author 2YSP
* @date 2020/8/19 20:06
*/
public class SendHandlerV2 extends ChannelInboundHandlerAdapter {
private static Logger logger = LoggerFactory.getLogger(SendHandlerV2.class);
/**
* 等待通道建立最大时间
*/
static final int CHANNEL_WAIT_TIME = 4;
/**
* 等待响应最大时间
*/
static final int RESPONSE_WAIT_TIME = 8;
private volatile Channel channel;
private String remoteAddress;
private static Map<String, RpcFuture<RpcResponse>> requestMap = new ConcurrentHashMap<>();
private MessageProtocol messageProtocol;
private CountDownLatch latch = new CountDownLatch(1);
public SendHandlerV2(MessageProtocol messageProtocol,String remoteAddress) {
this.messageProtocol = messageProtocol;
this.remoteAddress = remoteAddress;
}
@Override
public void channelRegistered(ChannelHandlerContext ctx) throws Exception {
this.channel = ctx.channel();
latch.countDown();
}
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
logger.debug("Connect to server successfully:{}", ctx);
}
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
logger.debug("Client reads message:{}", msg);
ByteBuf byteBuf = (ByteBuf) msg;
byte[] resp = new byte[byteBuf.readableBytes()];
byteBuf.readBytes(resp);
// 手动回收
ReferenceCountUtil.release(byteBuf);
RpcResponse response = messageProtocol.unmarshallingResponse(resp);
RpcFuture<RpcResponse> future = requestMap.get(response.getRequestId());
future.setResponse(response);
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
cause.printStackTrace();
logger.error("Exception occurred:{}", cause.getMessage());
ctx.close();
}
@Override
public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
ctx.flush();
}
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
super.channelInactive(ctx);
logger.error("channel inactive with remoteAddress:[{}]",remoteAddress);
NettyNetClient.connectedServerNodes.remove(remoteAddress);
}
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception {
super.userEventTriggered(ctx, evt);
}
public RpcResponse sendRequest(RpcRequest request) {
RpcResponse response;
RpcFuture<RpcResponse> future = new RpcFuture<>();
requestMap.put(request.getRequestId(), future);
try {
byte[] data = messageProtocol.marshallingRequest(request);
ByteBuf reqBuf = Unpooled.buffer(data.length);
reqBuf.writeBytes(data);
if (latch.await(CHANNEL_WAIT_TIME,TimeUnit.SECONDS)){
channel.writeAndFlush(reqBuf);
// 等待响应
response = future.get(RESPONSE_WAIT_TIME, TimeUnit.SECONDS);
}else {
throw new RpcException("establish channel time out");
}
} catch (Exception e) {
throw new RpcException(e.getMessage());
} finally {
requestMap.remove(request.getRequestId());
}
return response;
}
}
RpcFuture
package cn.sp.rpc.client.net;
import java.util.concurrent.*;
/**
* @author 2YSP
* @date 2020/8/19 22:31
*/
public class RpcFuture<T> implements Future<T> {
private T response;
/**
* 因为请求和响应是一一对应的,所以这里是1
*/
private CountDownLatch countDownLatch = new CountDownLatch(1);
/**
* Future的请求时间,用于计算Future是否超时
*/
private long beginTime = System.currentTimeMillis();
@Override
public boolean cancel(boolean mayInterruptIfRunning) {
return false;
}
@Override
public boolean isCancelled() {
return false;
}
@Override
public boolean isDone() {
if (response != null) {
return true;
}
return false;
}
/**
* 获取响应,直到有结果才返回
* @return
* @throws InterruptedException
* @throws ExecutionException
*/
@Override
public T get() throws InterruptedException, ExecutionException {
countDownLatch.await();
return response;
}
@Override
public T get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
if (countDownLatch.await(timeout,unit)){
return response;
}
return null;
}
public void setResponse(T response) {
this.response = response;
countDownLatch.countDown();
}
public long getBeginTime() {
return beginTime;
}
}
此处逻辑,第一次执行 SendHandlerV2#sendRequest() 时channel需要等待通道建立好之后才能发送请求,所以用CountDownLatch来控制,等待通道建立。
自定义Future+requestMap缓存来实现netty的请求和阻塞等待响应,RpcRequest对象在创建时会生成一个请求的唯一标识requestId,发送请求前先将RpcFuture缓存到requestMap中,key为requestId,读取到服务端的响应信息后(channelRead方法),将响应结果放入对应的RpcFuture中。
SendHandlerV2#channelInactive() 方法中,如果连接的服务端异常断开连接了,则及时清理缓存中对应的serverNode。
测试环境:
1.本地启动zookeeper
2.本地启动一个消费者,两个服务端,轮询算法
3.使用ab进行压力测试,4个线程发送10000个请求
ab -c 4 -n 10000 http://localhost:8080/test/user?id=1
测试结果: 从图片可以看出,10000个请求只用了11s,比之前的130+秒耗时减少了10倍以上。
代码地址:
本文由哈喽比特于3年以前收录,如有侵权请联系我们。
文章来源:https://mp.weixin.qq.com/s/fc_1aaM2pfZLPnkoH_5h9g
京东创始人刘强东和其妻子章泽天最近成为了互联网舆论关注的焦点。有关他们“移民美国”和在美国购买豪宅的传言在互联网上广泛传播。然而,京东官方通过微博发言人发布的消息澄清了这些传言,称这些言论纯属虚假信息和蓄意捏造。
日前,据博主“@超能数码君老周”爆料,国内三大运营商中国移动、中国电信和中国联通预计将集体采购百万台规模的华为Mate60系列手机。
据报道,荷兰半导体设备公司ASML正看到美国对华遏制政策的负面影响。阿斯麦(ASML)CEO彼得·温宁克在一档电视节目中分享了他对中国大陆问题以及该公司面临的出口管制和保护主义的看法。彼得曾在多个场合表达了他对出口管制以及中荷经济关系的担忧。
今年早些时候,抖音悄然上线了一款名为“青桃”的 App,Slogan 为“看见你的热爱”,根据应用介绍可知,“青桃”是一个属于年轻人的兴趣知识视频平台,由抖音官方出品的中长视频关联版本,整体风格有些类似B站。
日前,威马汽车首席数据官梅松林转发了一份“世界各国地区拥车率排行榜”,同时,他发文表示:中国汽车普及率低于非洲国家尼日利亚,每百户家庭仅17户有车。意大利世界排名第一,每十户中九户有车。
近日,一项新的研究发现,维生素 C 和 E 等抗氧化剂会激活一种机制,刺激癌症肿瘤中新血管的生长,帮助它们生长和扩散。
据媒体援引消息人士报道,苹果公司正在测试使用3D打印技术来生产其智能手表的钢质底盘。消息传出后,3D系统一度大涨超10%,不过截至周三收盘,该股涨幅回落至2%以内。
9月2日,坐拥千万粉丝的网红主播“秀才”账号被封禁,在社交媒体平台上引发热议。平台相关负责人表示,“秀才”账号违反平台相关规定,已封禁。据知情人士透露,秀才近期被举报存在违法行为,这可能是他被封禁的部分原因。据悉,“秀才”年龄39岁,是安徽省亳州市蒙城县人,抖音网红,粉丝数量超1200万。他曾被称为“中老年...
9月3日消息,亚马逊的一些股东,包括持有该公司股票的一家养老基金,日前对亚马逊、其创始人贝索斯和其董事会提起诉讼,指控他们在为 Project Kuiper 卫星星座项目购买发射服务时“违反了信义义务”。
据消息,为推广自家应用,苹果现推出了一个名为“Apps by Apple”的网站,展示了苹果为旗下产品(如 iPhone、iPad、Apple Watch、Mac 和 Apple TV)开发的各种应用程序。
特斯拉本周在美国大幅下调Model S和X售价,引发了该公司一些最坚定支持者的不满。知名特斯拉多头、未来基金(Future Fund)管理合伙人加里·布莱克发帖称,降价是一种“短期麻醉剂”,会让潜在客户等待进一步降价。
据外媒9月2日报道,荷兰半导体设备制造商阿斯麦称,尽管荷兰政府颁布的半导体设备出口管制新规9月正式生效,但该公司已获得在2023年底以前向中国运送受限制芯片制造机器的许可。
近日,根据美国证券交易委员会的文件显示,苹果卫星服务提供商 Globalstar 近期向马斯克旗下的 SpaceX 支付 6400 万美元(约 4.65 亿元人民币)。用于在 2023-2025 年期间,发射卫星,进一步扩展苹果 iPhone 系列的 SOS 卫星服务。
据报道,马斯克旗下社交平台𝕏(推特)日前调整了隐私政策,允许 𝕏 使用用户发布的信息来训练其人工智能(AI)模型。新的隐私政策将于 9 月 29 日生效。新政策规定,𝕏可能会使用所收集到的平台信息和公开可用的信息,来帮助训练 𝕏 的机器学习或人工智能模型。
9月2日,荣耀CEO赵明在采访中谈及华为手机回归时表示,替老同事们高兴,觉得手机行业,由于华为的回归,让竞争充满了更多的可能性和更多的魅力,对行业来说也是件好事。
《自然》30日发表的一篇论文报道了一个名为Swift的人工智能(AI)系统,该系统驾驶无人机的能力可在真实世界中一对一冠军赛里战胜人类对手。
近日,非营利组织纽约真菌学会(NYMS)发出警告,表示亚马逊为代表的电商平台上,充斥着各种AI生成的蘑菇觅食科普书籍,其中存在诸多错误。
社交媒体平台𝕏(原推特)新隐私政策提到:“在您同意的情况下,我们可能出于安全、安保和身份识别目的收集和使用您的生物识别信息。”
2023年德国柏林消费电子展上,各大企业都带来了最新的理念和产品,而高端化、本土化的中国产品正在不断吸引欧洲等国际市场的目光。
罗永浩日前在直播中吐槽苹果即将推出的 iPhone 新品,具体内容为:“以我对我‘子公司’的了解,我认为 iPhone 15 跟 iPhone 14 不会有什么区别的,除了序(列)号变了,这个‘不要脸’的东西,这个‘臭厨子’。