利用SpringBoot實現(xiàn)一個基于本地代理模式的RPC調(diào)用框架
在微服務架構(gòu)中,服務間通信是一個核心問題。
雖然Dubbo、gRPC等成熟框架已經(jīng)為我們提供了完整的RPC解決方案,但理解其底層原理并動手實現(xiàn)一個簡化版本,對提升我們的技術(shù)理解深度很有幫助。
本文將帶你從零開始,使用SpringBoot實現(xiàn)一個基于本地代理模式的RPC調(diào)用框架。
整體設計思路
我們的RPC框架采用經(jīng)典的代理模式設計:
接口定義:定義服務接口,客戶端和服務端共享
動態(tài)代理:客戶端通過JDK動態(tài)代理生成接口實現(xiàn)類
序列化:使用JSON進行數(shù)據(jù)序列化傳輸
網(wǎng)絡通信:基于HTTP協(xié)議進行服務間通信
服務注冊:服務端暴露接口實現(xiàn),客戶端動態(tài)發(fā)現(xiàn)
核心代碼實現(xiàn)
1. 定義RPC注解
首先創(chuàng)建用于標識RPC服務的注解:
// RpcService.java - 服務提供者注解
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Component
public @interface RpcService {
Class<?> value() default void.class;
String version() default "1.0";
}
// RpcReference.java - 服務消費者注解
@Target(ElementType.FIELD)
@Retention(RetentionPolicy.RUNTIME)
public @interface RpcReference {
String version() default "1.0";
long timeout() default 5000;
}
2. RPC請求響應模型
定義網(wǎng)絡傳輸?shù)臄?shù)據(jù)結(jié)構(gòu):
// RpcRequest.java
public class RpcRequest {
private String requestId;
private String className;
private String methodName;
private Class<?>[] parameterTypes;
private Object[] parameters;
private String version;
// 構(gòu)造函數(shù)
public RpcRequest() {
this.requestId = UUID.randomUUID().toString();
}
// getter/setter方法
public String getRequestId() { return requestId; }
public void setRequestId(String requestId) { this.requestId = requestId; }
public String getClassName() { return className; }
public void setClassName(String className) { this.className = className; }
public String getMethodName() { return methodName; }
public void setMethodName(String methodName) { this.methodName = methodName; }
public Class<?>[] getParameterTypes() { return parameterTypes; }
public void setParameterTypes(Class<?>[] parameterTypes) { this.parameterTypes = parameterTypes; }
public Object[] getParameters() { return parameters; }
public void setParameters(Object[] parameters) { this.parameters = parameters; }
public String getVersion() { return version; }
public void setVersion(String version) { this.version = version; }
}
// RpcResponse.java
public class RpcResponse {
private String requestId;
private Object result;
private String error;
private boolean success;
public RpcResponse() {}
public RpcResponse(String requestId) {
this.requestId = requestId;
}
// getter/setter方法
public String getRequestId() { return requestId; }
public void setRequestId(String requestId) { this.requestId = requestId; }
public Object getResult() { return result; }
public void setResult(Object result) {
this.result = result;
this.success = true;
}
public String getError() { return error; }
public void setError(String error) {
this.error = error;
this.success = false;
}
public boolean isSuccess() { return success; }
public void setSuccess(boolean success) { this.success = success; }
}
3. 服務注冊中心
實現(xiàn)簡單的本地服務注冊機制:
// ServiceRegistry.java
@Component
public class ServiceRegistry {
private static final Logger logger = LoggerFactory.getLogger(ServiceRegistry.class);
// 服務實例注冊表:接口名 -> 服務實現(xiàn)實例
private final Map<String, Object> serviceMap = new ConcurrentHashMap<>();
/**
* 注冊服務實例
*/
public void registerService(Class<?> serviceInterface, String version, Object serviceImpl) {
String serviceName = generateServiceName(serviceInterface, version);
serviceMap.put(serviceName, serviceImpl);
logger.info("注冊服務成功: {} -> {}", serviceName, serviceImpl.getClass().getName());
}
/**
* 獲取服務實例
*/
public Object getService(String className, String version) {
String serviceName = generateServiceName(className, version);
Object service = serviceMap.get(serviceName);
if (service == null) {
logger.warn("未找到服務: {}", serviceName);
}
return service;
}
/**
* 生成服務名稱
*/
private String generateServiceName(Class<?> serviceInterface, String version) {
return generateServiceName(serviceInterface.getName(), version);
}
private String generateServiceName(String className, String version) {
return className + ":" + version;
}
/**
* 獲取所有已注冊的服務
*/
public Set<String> getAllServices() {
return new HashSet<>(serviceMap.keySet());
}
}
4. RPC客戶端代理工廠
這是框架的核心,通過動態(tài)代理實現(xiàn)透明的遠程調(diào)用:
// RpcClientProxy.java
@Component
public class RpcClientProxy {
private static final Logger logger = LoggerFactory.getLogger(RpcClientProxy.class);
@Autowired
private RpcClient rpcClient;
/**
* 為指定接口創(chuàng)建代理實例
*/
@SuppressWarnings("unchecked")
public <T> T createProxy(Class<T> interfaceClass, String version, long timeout) {
return (T) Proxy.newProxyInstance(
interfaceClass.getClassLoader(),
new Class[]{interfaceClass},
new RpcInvocationHandler(interfaceClass, version, timeout, rpcClient)
);
}
/**
* 動態(tài)代理調(diào)用處理器
*/
private static class RpcInvocationHandler implements InvocationHandler {
private final Class<?> interfaceClass;
private final String version;
private final long timeout;
private final RpcClient rpcClient;
public RpcInvocationHandler(Class<?> interfaceClass, String version, long timeout, RpcClient rpcClient) {
this.interfaceClass = interfaceClass;
this.version = version;
this.timeout = timeout;
this.rpcClient = rpcClient;
}
@Override
public Object invoke(Object proxy, Method method, Object[] args) throws Throwable {
// 跳過Object類的基礎方法
if (Object.class.equals(method.getDeclaringClass())) {
return method.invoke(this, args);
}
// 構(gòu)建RPC請求
RpcRequest request = buildRpcRequest(method, args);
try {
// 發(fā)送遠程調(diào)用請求
RpcResponse response = rpcClient.sendRequest(request, timeout);
if (response.isSuccess()) {
return response.getResult();
} else {
throw new RuntimeException("RPC調(diào)用失敗: " + response.getError());
}
} catch (Exception e) {
logger.error("RPC調(diào)用異常: {}.{}", interfaceClass.getName(), method.getName(), e);
throw new RuntimeException("RPC調(diào)用異常", e);
}
}
/**
* 構(gòu)建RPC請求對象
*/
private RpcRequest buildRpcRequest(Method method, Object[] args) {
RpcRequest request = new RpcRequest();
request.setClassName(interfaceClass.getName());
request.setMethodName(method.getName());
request.setParameterTypes(method.getParameterTypes());
request.setParameters(args);
request.setVersion(version);
return request;
}
}
}
5. RPC網(wǎng)絡客戶端
負責實際的網(wǎng)絡通信:
// RpcClient.java
@Component
public class RpcClient {
private static final Logger logger = LoggerFactory.getLogger(RpcClient.class);
@Autowired
private RestTemplate restTemplate;
@Value("${rpc.server.url:http://localhost:8080}")
private String serverUrl;
/**
* 發(fā)送RPC請求
*/
public RpcResponse sendRequest(RpcRequest request, long timeout) {
try {
logger.debug("發(fā)送RPC請求: {}.{}", request.getClassName(), request.getMethodName());
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_JSON);
HttpEntity<RpcRequest> entity = new HttpEntity<>(request, headers);
// 發(fā)送HTTP POST請求
ResponseEntity<RpcResponse> responseEntity = restTemplate.postForEntity(
serverUrl + "/rpc/invoke",
entity,
RpcResponse.class
);
RpcResponse response = responseEntity.getBody();
logger.debug("收到RPC響應: requestId={}, success={}",
response.getRequestId(), response.isSuccess());
return response;
} catch (Exception e) {
logger.error("RPC請求發(fā)送失敗", e);
RpcResponse errorResponse = new RpcResponse(request.getRequestId());
errorResponse.setError("網(wǎng)絡請求失敗: " + e.getMessage());
return errorResponse;
}
}
}
6. RPC服務端處理器
處理客戶端發(fā)來的RPC請求:
@RestController
@RequestMapping("/rpc")
public class RpcRequestHandler {
private static final Logger logger = LoggerFactory.getLogger(RpcRequestHandler.class);
@Autowired
private ServiceRegistry serviceRegistry;
/**
* 處理RPC調(diào)用請求
*/
@PostMapping("/invoke")
public RpcResponse handleRpcRequest(@RequestBody RpcRequest request) {
RpcResponse response = new RpcResponse(request.getRequestId());
try {
logger.debug("處理RPC請求: {}.{}", request.getClassName(), request.getMethodName());
// 查找服務實例
Object serviceInstance = serviceRegistry.getService(request.getClassName(), request.getVersion());
if (serviceInstance == null) {
response.setError("服務未找到: " + request.getClassName());
return response;
}
// 通過反射調(diào)用方法
Class<?> serviceClass = serviceInstance.getClass();
Method method = serviceClass.getMethod(request.getMethodName(), request.getParameterTypes());
Object[] parameters = request.getParameters();
Class<?>[] paramTypes = method.getParameterTypes();
for (int i = 0; i < parameters.length; i++) {
if (parameters[i] != null && ClassUtil.isBasicType(paramTypes[i])) {
// 處理基本類型轉(zhuǎn)換(如客戶端傳的是包裝類型,服務端是基本類型)
parameters[i] = convertType(paramTypes[i], parameters[i]);
}
}
Object result = method.invoke(serviceInstance, parameters);
response.setResult(result);
logger.debug("RPC調(diào)用成功: {}.{}", request.getClassName(), request.getMethodName());
} catch (Exception e) {
logger.error("RPC調(diào)用處理異常", e);
response.setError("方法調(diào)用異常: " + e.getMessage());
}
return response;
}
/**
* 類型轉(zhuǎn)換處理(支持基本類型和包裝類互轉(zhuǎn))
*/
private Object convertType(Class<?> targetType, Object value) {
// 處理null值
if (value == null) return null;
// 類型匹配時直接返回
if (targetType.isInstance(value)) {
return value;
}
// 處理數(shù)字類型轉(zhuǎn)換
if (value instanceof Number) {
Number number = (Number) value;
if (targetType == int.class || targetType == Integer.class) return number.intValue();
if (targetType == long.class || targetType == Long.class) return number.longValue();
if (targetType == double.class || targetType == Double.class) return number.doubleValue();
if (targetType == float.class || targetType == Float.class) return number.floatValue();
if (targetType == byte.class || targetType == Byte.class) return number.byteValue();
if (targetType == short.class || targetType == Short.class) return number.shortValue();
}
// 處理布爾類型轉(zhuǎn)換
if (targetType == boolean.class || targetType == Boolean.class) {
if (value instanceof Boolean) return value;
return Boolean.parseBoolean(value.toString());
}
// 處理字符類型轉(zhuǎn)換
if (targetType == char.class || targetType == Character.class) {
String str = value.toString();
if (!str.isEmpty()) return str.charAt(0);
}
throw new IllegalArgumentException(String.format(
"類型轉(zhuǎn)換失敗: %s -> %s",
value.getClass().getSimpleName(),
targetType.getSimpleName()
));
}
/**
* 查詢已注冊的服務列表
*/
@GetMapping("/services")
public Set<String> getRegisteredServices() {
return serviceRegistry.getAllServices();
}
}
7. 自動配置和Bean后處理器
實現(xiàn)Spring Boot的自動裝配:
// RpcAutoConfiguration.java
@Configuration
@ComponentScan(basePackages = "com.example.rpc")
public class RpcAutoConfiguration {
@Bean
@ConditionalOnMissingBean
public RestTemplate restTemplate() {
RestTemplate restTemplate = new RestTemplate();
// 設置連接超時
HttpComponentsClientHttpRequestFactory factory = new HttpComponentsClientHttpRequestFactory();
factory.setConnectTimeout(3000);
factory.setReadTimeout(10000);
restTemplate.setRequestFactory(factory);
return restTemplate;
}
}
// RpcServiceProcessor.java - 處理@RpcService注解
@Component
public class RpcServiceProcessor implements BeanPostProcessor, ApplicationContextAware {
private static final Logger logger = LoggerFactory.getLogger(RpcServiceProcessor.class);
private ApplicationContext applicationContext;
@Override
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
this.applicationContext = applicationContext;
}
@Override
public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
Class<?> beanClass = bean.getClass();
// 檢查是否有@RpcService注解
RpcService rpcService = beanClass.getAnnotation(RpcService.class);
if (rpcService != null) {
registerRpcService(bean, rpcService);
}
return bean;
}
/**
* 注冊RPC服務
*/
private void registerRpcService(Object serviceBean, RpcService rpcService) {
ServiceRegistry serviceRegistry = applicationContext.getBean(ServiceRegistry.class);
Class<?> interfaceClass = rpcService.value();
if (interfaceClass == void.class) {
// 如果沒有指定接口,自動查找第一個接口
Class<?>[] interfaces = serviceBean.getClass().getInterfaces();
if (interfaces.length > 0) {
interfaceClass = interfaces[0];
} else {
logger.warn("無法確定服務接口: {}", serviceBean.getClass().getName());
return;
}
}
serviceRegistry.registerService(interfaceClass, rpcService.version(), serviceBean);
}
}
// RpcReferenceProcessor.java - 處理@RpcReference注解
@Component
public class RpcReferenceProcessor implements BeanPostProcessor, ApplicationContextAware {
private static final Logger logger = LoggerFactory.getLogger(RpcReferenceProcessor.class);
private ApplicationContext applicationContext;
@Override
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
this.applicationContext = applicationContext;
}
@Override
public Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException {
Class<?> beanClass = bean.getClass();
// 處理所有標注了@RpcReference的字段
Field[] fields = beanClass.getDeclaredFields();
for (Field field : fields) {
RpcReference rpcReference = field.getAnnotation(RpcReference.class);
if (rpcReference != null) {
injectRpcReference(bean, field, rpcReference);
}
}
return bean;
}
/**
* 注入RPC服務代理
*/
private void injectRpcReference(Object bean, Field field, RpcReference rpcReference) {
try {
RpcClientProxy rpcClientProxy = applicationContext.getBean(RpcClientProxy.class);
Class<?> interfaceClass = field.getType();
Object proxyInstance = rpcClientProxy.createProxy(
interfaceClass,
rpcReference.version(),
rpcReference.timeout()
);
field.setAccessible(true);
field.set(bean, proxyInstance);
logger.info("注入RPC服務代理: {}", interfaceClass.getName());
} catch (Exception e) {
logger.error("RPC服務代理注入失敗: {}", field.getName(), e);
throw new RuntimeException("RPC服務代理注入失敗", e);
}
}
}
使用示例
1. 定義服務接口
// UserService.java
public interface UserService {
String getUserName(Long userId);
boolean updateUser(Long userId, String name);
List<String> getUserList();
}
2. 實現(xiàn)服務提供者
// UserServiceImpl.java
@RpcService(UserService.class)
public class UserServiceImpl implements UserService {
private static final Map<Long, String> userDatabase = new ConcurrentHashMap<>();
static {
userDatabase.put(1L, "張三");
userDatabase.put(2L, "李四");
userDatabase.put(3L, "王五");
}
@Override
public String getUserName(Long userId) {
String userName = userDatabase.get(userId);
return userName != null ? userName : "用戶不存在";
}
@Override
public boolean updateUser(Long userId, String name) {
if (userDatabase.containsKey(userId)) {
userDatabase.put(userId, name);
return true;
}
return false;
}
@Override
public List<String> getUserList() {
return new ArrayList<>(userDatabase.values());
}
}
3. 創(chuàng)建服務消費者
// UserController.java
@RestController
@RequestMapping("/user")
public class UserController {
@RpcReference
private UserService userService;
@GetMapping("/{userId}")
public String getUser(@PathVariable Long userId) {
return userService.getUserName(userId);
}
@PostMapping("/{userId}")
public boolean updateUser(@PathVariable Long userId, @RequestParam String name) {
return userService.updateUser(userId, name);
}
@GetMapping("/list")
public List<String> getUserList() {
return userService.getUserList();
}
}
4. 啟動類配置
// Application.java
@SpringBootApplication
@EnableAutoConfiguration
public class RpcDemoApplication {
public static void main(String[] args) {
SpringApplication.run(RpcDemoApplication.class, args);
}
}
5. 配置文件
# application.yml
server:
port: 8080
rpc:
server:
url: http://localhost:8080
logging:
level:
com.example.rpc: DEBUG
測試驗證
啟動應用后,可以通過以下方式測試:
# 查詢用戶信息 curl http://localhost:8080/user/1 # 更新用戶信息 curl -X POST "http://localhost:8080/user/1?name=新名字" # 獲取用戶列表 curl http://localhost:8080/user/list # 查看已注冊的服務 curl http://localhost:8080/rpc/services
總結(jié)
這個實現(xiàn)雖然相對簡單,但完整展現(xiàn)了RPC框架的核心思想。
在實際項目中,建議使用成熟的RPC框架如Dubbo或Spring Cloud,但理解底層原理對我們選擇和優(yōu)化技術(shù)方案很有價值。
到此這篇關(guān)于利用SpringBoot實現(xiàn)一個基于本地代理模式的RPC調(diào)用框架的文章就介紹到這了,更多相關(guān)SpringBoot實現(xiàn)RPC調(diào)用內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Java8 自定義CompletableFuture的原理解析
這篇文章主要介紹了Java8 自定義CompletableFuture的原理解析,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-11-11
SpringBoot整合Swagger頁面禁止訪問swagger-ui.html方式
本文介紹了如何在SpringBoot項目中通過配置SpringSecurity和創(chuàng)建攔截器來禁止訪問SwaggerUI頁面,此外,還提供了禁用SwaggerUI和Swagger資源的配置方法,以確保這些端點和頁面對外部用戶不可見或無法訪問2025-02-02
fastjson全局日期序列化設置導致JSONField失效問題解決方案
這篇文章主要介紹了fastjson通過代碼指定全局序列化返回時間格式,導致使用JSONField注解標注屬性的特殊日期返回格式失效問題的解決方案2023-01-01
jar包在windows后臺運行,通過.bat文件實現(xiàn)
文章主要講解了在Windows后臺運行JAR包的方法,通過創(chuàng)建啟動.bat和停止.bat文件,使用javaw方式運行JAR包,可以在關(guān)閉cmd界面后仍然保持程序運行,提供了啟動.bat和停止.bat的具體內(nèi)容2026-04-04

