Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,11 @@ jobs:
distribution: 'temurin'
cache: 'maven'

- name: Run core and transport unit tests
run: mvn -B -ntp test -pl rpc-core,rpc-transport-netty -am -Drpc.registry=local
- name: Run core, transport and starter unit tests
run: >-
mvn -B -ntp test
-pl rpc-core,rpc-transport-netty,rpc-spring-boot-starter -am
-Drpc.registry=local

- name: Upload unit test and coverage reports
if: always()
Expand Down
23 changes: 22 additions & 1 deletion rpc-core/src/main/java/com/xiaoyu/rpc/core/client/RpcClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -76,8 +76,29 @@ public CompletableFuture<Object> sendRequest(RpcRequest request, Class<?> return

@Override
public void close() {
if (closed.compareAndSet(false, true)) {
if (!closed.compareAndSet(false, true)) {
return;
}

RuntimeException failure = null;
try {
transportClient.close();
} catch (RuntimeException e) {
failure = e;
}

try {
serviceDiscovery.close();
} catch (Exception e) {
if (failure == null) {
failure = new RuntimeException("关闭 ServiceDiscovery 失败", e);
} else {
failure.addSuppressed(e);
}
}

if (failure != null) {
throw failure;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,15 +5,20 @@
import java.net.InetSocketAddress;

/**
* 服务发现接口
* 服务发现接口。
*/
@SPI
public interface ServiceDiscovery {
public interface ServiceDiscovery extends AutoCloseable {
/**
* 查找服务地址
* 查找服务地址。
*
* @param serviceName 服务名称
* @return 服务地址
*/
InetSocketAddress lookupService(String serviceName);

@Override
default void close() {
// 默认实现无外部资源需要释放。
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,34 +2,31 @@

import com.alibaba.nacos.api.exception.NacosException;
import com.alibaba.nacos.api.naming.NamingService;
import com.alibaba.nacos.api.naming.pojo.Instance;
import com.alibaba.nacos.api.naming.listener.EventListener;
import com.alibaba.nacos.api.naming.listener.Event;
import com.alibaba.nacos.api.naming.listener.EventListener;
import com.alibaba.nacos.api.naming.listener.NamingEvent;
import com.alibaba.nacos.api.naming.pojo.Instance;
import com.xiaoyu.rpc.common.extension.ExtensionLoader;
import com.xiaoyu.rpc.core.config.RpcConfig;
import com.xiaoyu.rpc.core.loadbalancer.LoadBalancer;
import com.xiaoyu.rpc.core.registry.ServiceDiscovery;
import lombok.extern.slf4j.Slf4j;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.net.InetSocketAddress;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
import java.util.concurrent.atomic.AtomicBoolean;

public class NacosServiceDiscovery implements ServiceDiscovery {
private static final Logger log = LoggerFactory.getLogger(NacosServiceDiscovery.class);

private final NamingService namingService;
private final LoadBalancer loadBalancer;
// 本地缓存,用于容错和防抖
private static final java.util.Map<String, List<Instance>> serviceCache = new java.util.concurrent.ConcurrentHashMap<>();
// 已订阅的服务集合
private static final java.util.Set<String> subscribedServices = java.util.concurrent.ConcurrentHashMap.newKeySet();
private final Map<String, List<Instance>> serviceCache = new ConcurrentHashMap<>();
private final Map<String, EventListener> subscriptions = new ConcurrentHashMap<>();
private final AtomicBoolean closed = new AtomicBoolean(false);

public NacosServiceDiscovery() {
this(NacosUtils.getNacosNamingService(),
Expand All @@ -44,73 +41,96 @@ public NacosServiceDiscovery() {

@Override
public InetSocketAddress lookupService(String serviceName) {
if (closed.get()) {
throw new IllegalStateException("NacosServiceDiscovery 已关闭");
}

try {
// 第一次查找时订阅服务变更
if (subscribedServices.add(serviceName)) {
// add 返回 true 说明此前未订阅,避免同一个服务被重复订阅
subscribeService(serviceName);
}
ensureSubscribed(serviceName);

// 优先从 Nacos 拉取最新实例列表
List<Instance> instances = namingService.getAllInstances(serviceName);

if (instances.isEmpty()) {
log.warn("Nacos 返回实例列表为空,尝试使用本地缓存: {}", serviceName);
instances = serviceCache.get(serviceName);
} else {
// 更新本地缓存
serviceCache.put(serviceName, instances);
}

if (instances == null || instances.isEmpty()) {
log.error("未找到服务且本地无缓存: {}", serviceName);
// 注册中心明确返回空实例,代表当前服务已下线;不能继续使用旧缓存。
serviceCache.remove(serviceName);
throw new RuntimeException("未找到服务: " + serviceName);
}

// 转换 Instance 列表为 String 列表 (ip:port)
List<String> addressList = instances.stream()
.map(instance -> instance.getIp() + ":" + instance.getPort())
.collect(java.util.stream.Collectors.toList());

// 负载均衡选择
String targetAddress = loadBalancer.select(addressList);
log.info("负载均衡选择服务地址: {}", targetAddress);

String[] array = targetAddress.split(":");
return new InetSocketAddress(array[0], Integer.parseInt(array[1]));

updateCache(serviceName, instances);
return selectAddress(instances);
} catch (NacosException e) {
log.error("获取服务实例时发生网络异常,尝试回滚到本地缓存:", e);
// Nacos 短暂不可用时,优先用最近一次成功拉取到的实例兜底
// 只有 Nacos 网络/协议异常时,才允许使用最后一次成功结果容错。
log.error("获取服务实例时发生 Nacos 异常,尝试使用最近一次成功缓存: {}", serviceName, e);
List<Instance> cachedInstances = serviceCache.get(serviceName);
if (cachedInstances != null && !cachedInstances.isEmpty()) {
List<String> addressList = cachedInstances.stream()
.map(instance -> instance.getIp() + ":" + instance.getPort())
.collect(java.util.stream.Collectors.toList());
String targetAddress = loadBalancer.select(addressList);
String[] array = targetAddress.split(":");
return new InetSocketAddress(array[0], Integer.parseInt(array[1]));
return selectAddress(cachedInstances);
}
throw new RuntimeException("服务发现失败且无缓存可用: " + serviceName, e);
}
}

/**
* 订阅服务变更,实现本地缓存的实时更新
*/
private void subscribeService(String serviceName) throws NacosException {
namingService.subscribe(serviceName, new EventListener() {
@Override
public void onEvent(Event event) {
if (event instanceof NamingEvent) {
NamingEvent namingEvent = (NamingEvent) event;
List<Instance> instances = namingEvent.getInstances();
log.info("监听到服务变更,更新本地缓存: {} -> 实例数 {}", serviceName, instances.size());
if (instances != null && !instances.isEmpty()) {
serviceCache.put(serviceName, instances);
private void ensureSubscribed(String serviceName) throws NacosException {
if (subscriptions.containsKey(serviceName)) {
return;
}

synchronized (subscriptions) {
if (subscriptions.containsKey(serviceName)) {
return;
}

EventListener listener = new EventListener() {
@Override
public void onEvent(Event event) {
if (event instanceof NamingEvent) {
List<Instance> instances = ((NamingEvent) event).getInstances();
updateCache(serviceName, instances);
log.info("监听到服务变更,更新本地缓存: {} -> 实例数 {}",
serviceName, instances == null ? 0 : instances.size());
}
}
};
namingService.subscribe(serviceName, listener);
subscriptions.put(serviceName, listener);
}
}

void updateCache(String serviceName, List<Instance> instances) {
if (instances == null || instances.isEmpty()) {
serviceCache.remove(serviceName);
} else {
serviceCache.put(serviceName, List.copyOf(instances));
}
}

private InetSocketAddress selectAddress(List<Instance> instances) {
List<String> addressList = instances.stream()
.map(instance -> instance.getIp() + ":" + instance.getPort())
.toList();
String targetAddress = loadBalancer.select(addressList);
String[] array = targetAddress.split(":");
return new InetSocketAddress(array[0], Integer.parseInt(array[1]));
}

@Override
public void close() {
if (!closed.compareAndSet(false, true)) {
return;
}

subscriptions.forEach((serviceName, listener) -> {
try {
namingService.unsubscribe(serviceName, listener);
} catch (NacosException e) {
log.warn("取消 Nacos 服务订阅失败: {}", serviceName, e);
}
});
subscriptions.clear();
serviceCache.clear();

try {
namingService.shutDown();
} catch (NacosException e) {
log.warn("关闭 Nacos NamingService 失败", e);
}
}
}
38 changes: 32 additions & 6 deletions rpc-core/src/main/java/com/xiaoyu/rpc/core/server/RpcServer.java
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
import lombok.extern.slf4j.Slf4j;

import java.net.InetSocketAddress;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;

@Slf4j
Expand All @@ -17,6 +19,8 @@ public class RpcServer implements AutoCloseable {
private final int serverPort;
private final ServiceRegistry serviceRegistry;
private final TransportServer transportServer;
private final Set<String> registeredServices = ConcurrentHashMap.newKeySet();
private final AtomicBoolean started = new AtomicBoolean(false);
private final AtomicBoolean stopped = new AtomicBoolean(false);
private final Thread shutdownHook;

Expand Down Expand Up @@ -67,32 +71,49 @@ public <T> void register(Class<T> interfaceClass, T serviceImpl) {

String serviceName = interfaceClass.getName();
ServiceRepository.registerService(serviceName, serviceImpl);
registeredServices.add(serviceName);

try {
serviceRegistry.registerService(serviceName, new InetSocketAddress(serverHost, serverPort));
log.info("Service registered: {}", serviceName);
} catch (Exception e) {
log.error("Failed to register service: {}", serviceName, e);
if (started.get()) {
publishService(serviceName);
}
}

/**
* 启动传输层,确认端口已 bind 后再将服务发布到注册中心。
*/
public void start() throws InterruptedException {
if (stopped.get()) {
throw new IllegalStateException("RpcServer 已关闭");
}
if (!started.compareAndSet(false, true)) {
throw new IllegalStateException("RpcServer 已启动");
}

try {
transportServer.start();
} finally {
for (String serviceName : registeredServices) {
publishService(serviceName);
}
log.info("RPC Server 启动完成,已发布 {} 个服务", registeredServices.size());
} catch (InterruptedException | RuntimeException e) {
stop();
throw e;
}
}

/**
* 阻塞等待底层 TransportServer 停止。
*/
public void awaitTermination() throws InterruptedException {
transportServer.awaitTermination();
}

public void stop() {
if (!stopped.compareAndSet(false, true)) {
return;
}

// 先从注册中心摘除,阻止新流量,再停止网络层。
try {
serviceRegistry.clearRegistry();
} catch (Exception e) {
Expand All @@ -113,6 +134,11 @@ public void close() {
stop();
}

private void publishService(String serviceName) {
serviceRegistry.registerService(serviceName, new InetSocketAddress(serverHost, serverPort));
log.info("Service published after transport ready: {} -> {}:{}", serviceName, serverHost, serverPort);
}

private void removeShutdownHook() {
if (shutdownHook == null || Thread.currentThread() == shutdownHook) {
return;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,19 +1,28 @@
package com.xiaoyu.rpc.core.transport;

/**
* 传输层服务端接口
* 传输层服务端接口。
*/
public interface TransportServer {

/**
* 启动服务
*
* @throws InterruptedException 如果启动过程被中断
* 启动服务并在监听端口成功 bind 后返回。
*
* @throws InterruptedException 启动过程被中断
*/
void start() throws InterruptedException;

/**
* 停止服务
* 阻塞等待服务端停止。默认实现用于不需要阻塞语义的传输层。
*
* @throws InterruptedException 等待过程被中断
*/
default void awaitTermination() throws InterruptedException {
// no-op by default
}

/**
* 停止服务并释放资源。
*/
void stop();
}
Loading