Skip to content

Commit 4e8caab

Browse files
authored
Merge pull request #10 from xiaoyumuxi/fix/p1-service-lifecycle-discovery
fix: harden service startup and discovery lifecycle
2 parents cca65ea + 6e9bd13 commit 4e8caab

13 files changed

Lines changed: 472 additions & 195 deletions

File tree

‎.github/workflows/ci.yml‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -49,8 +49,11 @@ jobs:
4949
distribution: 'temurin'
5050
cache: 'maven'
5151

52-
- name: Run core and transport unit tests
53-
run: mvn -B -ntp test -pl rpc-core,rpc-transport-netty -am -Drpc.registry=local
52+
- name: Run core, transport and starter unit tests
53+
run: >-
54+
mvn -B -ntp test
55+
-pl rpc-core,rpc-transport-netty,rpc-spring-boot-starter -am
56+
-Drpc.registry=local
5457
5558
- name: Upload unit test and coverage reports
5659
if: always()

‎rpc-core/src/main/java/com/xiaoyu/rpc/core/client/RpcClient.java‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -76,8 +76,29 @@ public CompletableFuture<Object> sendRequest(RpcRequest request, Class<?> return
7676

7777
@Override
7878
public void close() {
79-
if (closed.compareAndSet(false, true)) {
79+
if (!closed.compareAndSet(false, true)) {
80+
return;
81+
}
82+
83+
RuntimeException failure = null;
84+
try {
8085
transportClient.close();
86+
} catch (RuntimeException e) {
87+
failure = e;
88+
}
89+
90+
try {
91+
serviceDiscovery.close();
92+
} catch (Exception e) {
93+
if (failure == null) {
94+
failure = new RuntimeException("关闭 ServiceDiscovery 失败", e);
95+
} else {
96+
failure.addSuppressed(e);
97+
}
98+
}
99+
100+
if (failure != null) {
101+
throw failure;
81102
}
82103
}
83104
}

‎rpc-core/src/main/java/com/xiaoyu/rpc/core/registry/ServiceDiscovery.java‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,15 +5,20 @@
55
import java.net.InetSocketAddress;
66

77
/**
8-
* 服务发现接口
8+
* 服务发现接口。
99
*/
1010
@SPI
11-
public interface ServiceDiscovery {
11+
public interface ServiceDiscovery extends AutoCloseable {
1212
/**
13-
* 查找服务地址
13+
* 查找服务地址。
1414
*
1515
* @param serviceName 服务名称
1616
* @return 服务地址
1717
*/
1818
InetSocketAddress lookupService(String serviceName);
19+
20+
@Override
21+
default void close() {
22+
// 默认实现无外部资源需要释放。
23+
}
1924
}

‎rpc-core/src/main/java/com/xiaoyu/rpc/core/registry/nacos/NacosServiceDiscovery.java‎

Lines changed: 78 additions & 58 deletions
Original file line numberDiff line numberDiff line change
@@ -2,34 +2,31 @@
22

33
import com.alibaba.nacos.api.exception.NacosException;
44
import com.alibaba.nacos.api.naming.NamingService;
5-
import com.alibaba.nacos.api.naming.pojo.Instance;
6-
import com.alibaba.nacos.api.naming.listener.EventListener;
75
import com.alibaba.nacos.api.naming.listener.Event;
6+
import com.alibaba.nacos.api.naming.listener.EventListener;
87
import com.alibaba.nacos.api.naming.listener.NamingEvent;
8+
import com.alibaba.nacos.api.naming.pojo.Instance;
99
import com.xiaoyu.rpc.common.extension.ExtensionLoader;
1010
import com.xiaoyu.rpc.core.config.RpcConfig;
1111
import com.xiaoyu.rpc.core.loadbalancer.LoadBalancer;
1212
import com.xiaoyu.rpc.core.registry.ServiceDiscovery;
13-
import lombok.extern.slf4j.Slf4j;
1413
import org.slf4j.Logger;
1514
import org.slf4j.LoggerFactory;
1615

1716
import java.net.InetSocketAddress;
1817
import java.util.List;
1918
import java.util.Map;
20-
import java.util.Set;
2119
import java.util.concurrent.ConcurrentHashMap;
22-
import java.util.stream.Collectors;
20+
import java.util.concurrent.atomic.AtomicBoolean;
2321

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

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

3431
public NacosServiceDiscovery() {
3532
this(NacosUtils.getNacosNamingService(),
@@ -44,73 +41,96 @@ public NacosServiceDiscovery() {
4441

4542
@Override
4643
public InetSocketAddress lookupService(String serviceName) {
44+
if (closed.get()) {
45+
throw new IllegalStateException("NacosServiceDiscovery 已关闭");
46+
}
47+
4748
try {
48-
// 第一次查找时订阅服务变更
49-
if (subscribedServices.add(serviceName)) {
50-
// add 返回 true 说明此前未订阅,避免同一个服务被重复订阅
51-
subscribeService(serviceName);
52-
}
49+
ensureSubscribed(serviceName);
5350

54-
// 优先从 Nacos 拉取最新实例列表
5551
List<Instance> instances = namingService.getAllInstances(serviceName);
56-
57-
if (instances.isEmpty()) {
58-
log.warn("Nacos 返回实例列表为空,尝试使用本地缓存: {}", serviceName);
59-
instances = serviceCache.get(serviceName);
60-
} else {
61-
// 更新本地缓存
62-
serviceCache.put(serviceName, instances);
63-
}
64-
6552
if (instances == null || instances.isEmpty()) {
66-
log.error("未找到服务且本地无缓存: {}", serviceName);
53+
// 注册中心明确返回空实例,代表当前服务已下线;不能继续使用旧缓存。
54+
serviceCache.remove(serviceName);
6755
throw new RuntimeException("未找到服务: " + serviceName);
6856
}
6957

70-
// 转换 Instance 列表为 String 列表 (ip:port)
71-
List<String> addressList = instances.stream()
72-
.map(instance -> instance.getIp() + ":" + instance.getPort())
73-
.collect(java.util.stream.Collectors.toList());
74-
75-
// 负载均衡选择
76-
String targetAddress = loadBalancer.select(addressList);
77-
log.info("负载均衡选择服务地址: {}", targetAddress);
78-
79-
String[] array = targetAddress.split(":");
80-
return new InetSocketAddress(array[0], Integer.parseInt(array[1]));
81-
58+
updateCache(serviceName, instances);
59+
return selectAddress(instances);
8260
} catch (NacosException e) {
83-
log.error("获取服务实例时发生网络异常,尝试回滚到本地缓存:", e);
84-
// Nacos 短暂不可用时,优先用最近一次成功拉取到的实例兜底
61+
// 只有 Nacos 网络/协议异常时,才允许使用最后一次成功结果容错。
62+
log.error("获取服务实例时发生 Nacos 异常,尝试使用最近一次成功缓存: {}", serviceName, e);
8563
List<Instance> cachedInstances = serviceCache.get(serviceName);
8664
if (cachedInstances != null && !cachedInstances.isEmpty()) {
87-
List<String> addressList = cachedInstances.stream()
88-
.map(instance -> instance.getIp() + ":" + instance.getPort())
89-
.collect(java.util.stream.Collectors.toList());
90-
String targetAddress = loadBalancer.select(addressList);
91-
String[] array = targetAddress.split(":");
92-
return new InetSocketAddress(array[0], Integer.parseInt(array[1]));
65+
return selectAddress(cachedInstances);
9366
}
9467
throw new RuntimeException("服务发现失败且无缓存可用: " + serviceName, e);
9568
}
9669
}
9770

98-
/**
99-
* 订阅服务变更,实现本地缓存的实时更新
100-
*/
101-
private void subscribeService(String serviceName) throws NacosException {
102-
namingService.subscribe(serviceName, new EventListener() {
103-
@Override
104-
public void onEvent(Event event) {
105-
if (event instanceof NamingEvent) {
106-
NamingEvent namingEvent = (NamingEvent) event;
107-
List<Instance> instances = namingEvent.getInstances();
108-
log.info("监听到服务变更,更新本地缓存: {} -> 实例数 {}", serviceName, instances.size());
109-
if (instances != null && !instances.isEmpty()) {
110-
serviceCache.put(serviceName, instances);
71+
private void ensureSubscribed(String serviceName) throws NacosException {
72+
if (subscriptions.containsKey(serviceName)) {
73+
return;
74+
}
75+
76+
synchronized (subscriptions) {
77+
if (subscriptions.containsKey(serviceName)) {
78+
return;
79+
}
80+
81+
EventListener listener = new EventListener() {
82+
@Override
83+
public void onEvent(Event event) {
84+
if (event instanceof NamingEvent) {
85+
List<Instance> instances = ((NamingEvent) event).getInstances();
86+
updateCache(serviceName, instances);
87+
log.info("监听到服务变更,更新本地缓存: {} -> 实例数 {}",
88+
serviceName, instances == null ? 0 : instances.size());
11189
}
11290
}
91+
};
92+
namingService.subscribe(serviceName, listener);
93+
subscriptions.put(serviceName, listener);
94+
}
95+
}
96+
97+
void updateCache(String serviceName, List<Instance> instances) {
98+
if (instances == null || instances.isEmpty()) {
99+
serviceCache.remove(serviceName);
100+
} else {
101+
serviceCache.put(serviceName, List.copyOf(instances));
102+
}
103+
}
104+
105+
private InetSocketAddress selectAddress(List<Instance> instances) {
106+
List<String> addressList = instances.stream()
107+
.map(instance -> instance.getIp() + ":" + instance.getPort())
108+
.toList();
109+
String targetAddress = loadBalancer.select(addressList);
110+
String[] array = targetAddress.split(":");
111+
return new InetSocketAddress(array[0], Integer.parseInt(array[1]));
112+
}
113+
114+
@Override
115+
public void close() {
116+
if (!closed.compareAndSet(false, true)) {
117+
return;
118+
}
119+
120+
subscriptions.forEach((serviceName, listener) -> {
121+
try {
122+
namingService.unsubscribe(serviceName, listener);
123+
} catch (NacosException e) {
124+
log.warn("取消 Nacos 服务订阅失败: {}", serviceName, e);
113125
}
114126
});
127+
subscriptions.clear();
128+
serviceCache.clear();
129+
130+
try {
131+
namingService.shutDown();
132+
} catch (NacosException e) {
133+
log.warn("关闭 Nacos NamingService 失败", e);
134+
}
115135
}
116136
}

‎rpc-core/src/main/java/com/xiaoyu/rpc/core/server/RpcServer.java‎

Lines changed: 32 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88
import lombok.extern.slf4j.Slf4j;
99

1010
import java.net.InetSocketAddress;
11+
import java.util.Set;
12+
import java.util.concurrent.ConcurrentHashMap;
1113
import java.util.concurrent.atomic.AtomicBoolean;
1214

1315
@Slf4j
@@ -17,6 +19,8 @@ public class RpcServer implements AutoCloseable {
1719
private final int serverPort;
1820
private final ServiceRegistry serviceRegistry;
1921
private final TransportServer transportServer;
22+
private final Set<String> registeredServices = ConcurrentHashMap.newKeySet();
23+
private final AtomicBoolean started = new AtomicBoolean(false);
2024
private final AtomicBoolean stopped = new AtomicBoolean(false);
2125
private final Thread shutdownHook;
2226

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

6872
String serviceName = interfaceClass.getName();
6973
ServiceRepository.registerService(serviceName, serviceImpl);
74+
registeredServices.add(serviceName);
7075

71-
try {
72-
serviceRegistry.registerService(serviceName, new InetSocketAddress(serverHost, serverPort));
73-
log.info("Service registered: {}", serviceName);
74-
} catch (Exception e) {
75-
log.error("Failed to register service: {}", serviceName, e);
76+
if (started.get()) {
77+
publishService(serviceName);
7678
}
7779
}
7880

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

8492
try {
8593
transportServer.start();
86-
} finally {
94+
for (String serviceName : registeredServices) {
95+
publishService(serviceName);
96+
}
97+
log.info("RPC Server 启动完成,已发布 {} 个服务", registeredServices.size());
98+
} catch (InterruptedException | RuntimeException e) {
8799
stop();
100+
throw e;
88101
}
89102
}
90103

104+
/**
105+
* 阻塞等待底层 TransportServer 停止。
106+
*/
107+
public void awaitTermination() throws InterruptedException {
108+
transportServer.awaitTermination();
109+
}
110+
91111
public void stop() {
92112
if (!stopped.compareAndSet(false, true)) {
93113
return;
94114
}
95115

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

137+
private void publishService(String serviceName) {
138+
serviceRegistry.registerService(serviceName, new InetSocketAddress(serverHost, serverPort));
139+
log.info("Service published after transport ready: {} -> {}:{}", serviceName, serverHost, serverPort);
140+
}
141+
116142
private void removeShutdownHook() {
117143
if (shutdownHook == null || Thread.currentThread() == shutdownHook) {
118144
return;
Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,28 @@
11
package com.xiaoyu.rpc.core.transport;
22

33
/**
4-
* 传输层服务端接口
4+
* 传输层服务端接口。
55
*/
66
public interface TransportServer {
77

88
/**
9-
* 启动服务
10-
*
11-
* @throws InterruptedException 如果启动过程被中断
9+
* 启动服务并在监听端口成功 bind 后返回。
10+
*
11+
* @throws InterruptedException 启动过程被中断
1212
*/
1313
void start() throws InterruptedException;
1414

1515
/**
16-
* 停止服务
16+
* 阻塞等待服务端停止。默认实现用于不需要阻塞语义的传输层。
17+
*
18+
* @throws InterruptedException 等待过程被中断
19+
*/
20+
default void awaitTermination() throws InterruptedException {
21+
// no-op by default
22+
}
23+
24+
/**
25+
* 停止服务并释放资源。
1726
*/
1827
void stop();
1928
}

0 commit comments

Comments
 (0)