智算多多
官方邮箱:sw@zsdodo.com

公司地址:北京市丰台区南四环西路188号总部基地三区国联股份数字经济总部(邮 编:100070)


京公网安备11010602202532号 本文深入剖析RPC核心本质与通用架构,详解Dubbo 3.x(Java生态企业级框架)和gRPC(云原生跨语言框架)的底层原理、性能差异、生产调优及避坑指南,涵盖动态代理、序列化、网络传输、服务发现、集群容错等关键模块,助力构建高可用分布式系统。
RPC(Remote Procedure Call,远程过程调用)的核心价值,是屏蔽分布式系统中网络通信、序列化、服务寻址、集群容错等底层复杂度,让开发者可以像调用本地方法一样,调用部署在远程节点的服务能力,是分布式微服务架构的核心基础设施。
一次完整的RPC调用,本质是一次跨进程的请求-响应交互,核心链路可拆解为6个核心阶段,流程如下:
所有成熟的RPC框架,都遵循分层设计理念,核心分为5层,每层职责单一且可扩展:
| 分层 | 核心职责 | 核心扩展点 |
|---|---|---|
| 服务接口层 | 定义服务对外暴露的接口与方法契约,屏蔽底层实现 | 服务注册、版本控制、分组隔离 |
| 代理层 | 为服务接口生成动态代理,封装远程调用的所有细节 | 动态代理实现(JDK/字节码生成) |
| 序列化层 | 实现对象与二进制流的双向转换,解决跨进程数据传输问题 | 序列化协议(Hessian2/Protobuf/Kryo等) |
| 网络传输层 | 实现二进制数据的跨网络可靠传输,处理网络连接、IO读写 | 传输协议(TCP/HTTP2)、IO模型(BIO/NIO/AIO) |
| 集群治理层 | 解决分布式环境下的服务发现、负载均衡、容错降级、流量管控问题 | 注册中心、负载均衡策略、集群容错机制 |
Dubbo是阿里开源的Java生态原生RPC框架,历经多年生产验证,3.x版本完成了云原生架构升级,是国内企业级微服务架构的主流选型。
代理层是RPC框架的入口,Dubbo的代理层核心是ProxyFactory接口,提供两种动态代理实现:
代理层的核心逻辑,是在消费端发起调用时,拦截所有方法调用,将方法名、参数类型、参数值、调用上下文封装为RpcInvocation请求对象,交给后续链路处理。
序列化是RPC性能的核心瓶颈之一,Dubbo提供了全场景的序列化协议适配,核心特性如下:
| 序列化协议 | 核心特点 | 适用场景 |
|---|---|---|
| Hessian2 | Dubbo默认协议,二进制序列化,跨语言兼容,性能稳定,支持对象循环引用 | 常规Java微服务业务场景,默认首选 |
| Protobuf | 谷歌开源的结构化数据序列化协议,压缩比极高,序列化速度快,强类型IDL,跨语言能力强 | 跨语言微服务、大报文传输、高性能要求场景 |
| Kryo/FST | Java专属序列化协议,性能远超Hessian2,压缩比更高 | Java生态内的高性能场景,大对象传输 |
| Fastjson2 | JSON格式序列化,可读性强,跨语言兼容 | 网关透传、调试场景、需要JSON明文的场景 |
序列化的核心性能指标有两个:序列化/反序列化耗时、序列化后的二进制体积,两者直接决定RPC调用的网络开销与CPU开销。
Dubbo 3.x 基于Netty 4.x 实现高性能NIO网络传输,原生支持TCP与HTTP2双协议,核心采用Reactor主从多线程模型:
为了避免业务逻辑阻塞IO线程,Dubbo设计了灵活的线程调度模型Dispatcher,核心实现如下:
Dubbo 3.x 完成了从接口级服务发现到应用级服务发现的架构升级,彻底对齐云原生Kubernetes的Service模型,核心优势如下:
服务发现的核心流程:服务提供者启动时,将自身的应用地址、端口、元数据信息注册到注册中心;消费者启动时,从注册中心订阅对应应用的地址列表,本地缓存并监听地址变化,实现服务地址的动态感知。
分布式环境下,服务调用不可避免会出现网络波动、节点宕机等问题,Dubbo提供了完善的集群容错机制,核心实现如下:
负载均衡是集群流量分发的核心,Dubbo提供了多种高性能负载均衡策略:
gRPC是Google开源的高性能、跨语言RPC框架,基于HTTP/2标准协议与Protocol Buffers(Protobuf)序列化协议设计,原生支持流式调用,是云原生、多语言微服务、跨平台服务调用的主流选型。
gRPC采用契约优先的设计理念,通过Protobuf IDL(接口定义语言)统一描述服务接口与数据结构,再通过protoc编译器与gRPC插件,生成对应语言的客户端与服务端代码,彻底解决跨语言的接口兼容问题。
Protobuf IDL的核心优势:
Protobuf是gRPC默认的唯一序列化协议,也是gRPC高性能的核心支撑,其底层采用TLV(Tag-Length-Value)存储结构与Varint变长编码,核心原理如下:
Varint变长编码
Varint是一种紧凑的数字编码方式,每个字节的最高位是标志位,1表示后续字节仍属于当前数字,0表示当前字节是数字的最后一个字节。例如:
对于绝大多数业务场景的数字,Varint编码可以将4字节的int压缩到1-2字节,8字节的long压缩到1-4字节,大幅降低序列化后的体积。
TLV存储结构
Protobuf的每个字段都由Tag、Length(可选)、Value三部分组成:
这种存储结构,使得Protobuf可以忽略不认识的字段,天然支持字段的新增与删除,实现向前向后兼容,同时序列化后的体积远小于JSON、XML等文本格式,序列化/反序列化速度是JSON的3-10倍。
gRPC底层完全基于HTTP/2协议设计,原生继承了HTTP/2的所有高性能特性,这也是gRPC与传统RPC框架最大的区别:
gRPC设计了可扩展的名称解析与负载均衡架构,核心分为两个组件:
gRPC的负载均衡是客户端侧实现的,客户端本地缓存服务端地址列表,直接发起点对点调用,无需经过中间代理,减少了网络跳转开销。
| 对比维度 | Dubbo 3.x | gRPC |
|---|---|---|
| 底层传输协议 | 原生支持TCP私有协议、HTTP/2协议,TCP协议性能更优 | 完全基于HTTP/2协议设计,协议通用性更强 |
| 序列化协议 | 多协议适配,默认Hessian2,支持Protobuf、Kryo、FST、JSON等 | 仅原生支持Protobuf,强绑定,序列化性能极致 |
| 服务治理能力 | 企业级全链路服务治理,内置注册中心、配置中心、流量管控、熔断降级、限流、链路追踪等全套能力 | 核心聚焦RPC调用本身,服务治理能力需通过拦截器、第三方组件扩展实现 |
| 跨语言能力 | Java生态原生,其他语言支持有限,跨语言能力较弱 | 原生支持几乎所有主流编程语言,跨语言能力极强 |
| 流式调用 | 3.x版本基于HTTP/2支持流式调用,能力完善度一般 | 原生深度支持4种流式调用模式,流式场景适配性极强 |
| 云原生适配 | 3.x版本完成云原生升级,支持应用级服务发现、K8s、Istio适配 | 云原生原生设计,与K8s、云原生网关、服务网格无缝适配,是CNCF毕业项目 |
| 性能表现 | Java生态内TCP协议场景下,性能优于gRPC;HTTP/2场景与gRPC持平 | 跨语言场景、流式场景、大报文场景下,性能优势明显 |
| 适用场景 | Java生态为主的企业级微服务架构,需要强服务治理能力的业务场景 | 多语言微服务、云原生架构、跨平台服务调用、流式数据传输场景 |
<dubbo:protocol name="dubbo" serialization="kryo"/><dubbo:protocol name="dubbo" iothreads="8"/><dubbo:protocol name="dubbo" payload="8388608" accepts="1000">
<dubbo:parameter key="tcp.nodelay" value="true"/>
<dubbo:parameter key="so.backlog" value="1024"/>
<dubbo:parameter key="so.keepalive" value="true"/>
</dubbo:protocol><dubbo:protocol name="dubbo" dispatcher="message"/><dubbo:protocol name="dubbo" threadpool="cached" threads="200" queues="0"/><dubbo:reference interface="com.jam.demo.service.UserService" retries="2"/><dubbo:application name="demo-provider" registry-mode="instance"/>ManagedChannel channel = ManagedChannelBuilder.forAddress("127.0.0.1", 9090)
.flowControlWindow(1024 * 1024)
.build();Server server = ServerBuilder.forPort(9090)
.maxConcurrentCallsPerConnection(1000)
.addService(new UserServiceImpl())
.build();ManagedChannel channel = ManagedChannelBuilder.forAddress("127.0.0.1", 9090)
.keepAliveTime(30, TimeUnit.SECONDS)
.keepAliveTimeout(5, TimeUnit.SECONDS)
.keepAliveWithoutCalls(true)
.build();Server server = ServerBuilder.forPort(9090)
.maxInboundMessageSize(16 * 1024 * 1024)
.addService(new UserServiceImpl())
.build();EventLoopGroup bossGroup = new NioEventLoopGroup(4);
EventLoopGroup workerGroup = new NioEventLoopGroup(8);
Server server = NettyServerBuilder.forPort(9090)
.bossEventLoopGroup(bossGroup)
.workerEventLoopGroup(workerGroup)
.addService(new UserServiceImpl())
.build();ExecutorService businessExecutor = new ThreadPoolExecutor(
40,
200,
60L,
TimeUnit.SECONDS,
new SynchronousQueue<>(),
new ThreadFactory() {
private final AtomicInteger threadNumber = new AtomicInteger(1);
@Override
public Thread newThread(Runnable r) {
return new Thread(r, "grpc-business-thread-" + threadNumber.getAndIncrement());
}
},
new ThreadPoolExecutor.AbortPolicy());
Server server = ServerBuilder.forPort(9090)
.executor(businessExecutor)
.addService(new UserServiceImpl())
.build();ManagedChannel channel = ManagedChannelBuilder.forAddress("127.0.0.1", 9090)
.defaultLoadBalancingPolicy("round_robin")
.build();<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.2.4</version>
<relativePath/>
</parent>
<groupId>com.jam.demo</groupId>
<artifactId>dubbo-demo</artifactId>
<version>1.0.0</version>
<name>dubbo-demo</name>
<properties>
<java.version>17</java.version>
<dubbo.version>3.2.10</dubbo.version>
<lombok.version>1.18.30</lombok.version>
<fastjson2.version>2.0.49</fastjson2.version>
<guava.version>33.1.0-jre</guava.version>
<springdoc.version>2.5.0</springdoc.version>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.dubbo</groupId>
<artifactId>dubbo-spring-boot-starter</artifactId>
<version>${dubbo.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springdoc</groupId>
<artifactId>springdoc-openapi-starter-webmvc-ui</artifactId>
<version>${springdoc.version}</version>
</dependency>
<dependency>
<groupId>org.apache.dubbo</groupId>
<artifactId>dubbo-registry-nacos</artifactId>
<version>${dubbo.version}</version>
</dependency>
<dependency>
<groupId>com.alibaba.nacos</groupId>
<artifactId>nacos-client</artifactId>
<version>2.3.2</version>
</dependency>
<dependency>
<groupId>com.esotericsoftware</groupId>
<artifactId>kryo</artifactId>
<version>5.5.0</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>${lombok.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>com.alibaba.fastjson2</groupId>
<artifactId>fastjson2</artifactId>
<version>${fastjson2.version}</version>
</dependency>
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
<version>${guava.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<excludes>
<exclude>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</exclude>
</excludes>
</configuration>
</plugin>
</plugins>
</build>
</project>
package com.jam.demo.service;
import com.jam.demo.dto.UserDTO;
import com.jam.demo.dto.UserQueryDTO;
import com.jam.demo.common.Result;
import java.util.List;
/**
* 用户服务RPC接口
* @author ken
*/
public interface UserService {
/**
* 根据用户ID查询用户信息
* @param userId 用户ID
* @return 用户信息
*/
Result<UserDTO> getUserById(Long userId);
/**
* 根据条件查询用户列表
* @param queryDTO 查询条件
* @return 用户列表
*/
Result<List<UserDTO>> listUserByCondition(UserQueryDTO queryDTO);
/**
* 新增用户信息
* @param userDTO 用户信息
* @return 新增结果
*/
Result<Long> addUser(UserDTO userDTO);
}
package com.jam.demo.dto;
import io.swagger.v3.oas.annotations.media.Schema;
import lombok.Data;
import java.io.Serializable;
import java.time.LocalDateTime;
/**
* 用户数据传输对象
* @author ken
*/
@Data
@Schema(description = "用户信息DTO")
public class UserDTO implements Serializable {
private static final long serialVersionUID = 1L;
@Schema(description = "用户ID", example = "1")
private Long userId;
@Schema(description = "用户名", example = "jam")
private String username;
@Schema(description = "用户昵称", example = "果酱")
private String nickname;
@Schema(description = "邮箱", example = "jam@demo.com")
private String email;
@Schema(description = "手机号", example = "13800138000")
private String phone;
@Schema(description = "创建时间")
private LocalDateTime createTime;
@Schema(description = "更新时间")
private LocalDateTime updateTime;
}
package com.jam.demo.service.impl;
import com.jam.demo.dto.UserDTO;
import com.jam.demo.dto.UserQueryDTO;
import com.jam.demo.common.Result;
import com.jam.demo.service.UserService;
import lombok.extern.slf4j.Slf4j;
import org.apache.dubbo.config.annotation.DubboService;
import org.springframework.util.StringUtils;
import org.springframework.util.CollectionUtils;
import com.google.common.collect.Lists;
import java.time.LocalDateTime;
import java.util.List;
import java.util.stream.Collectors;
/**
* 用户服务RPC实现类
* @author ken
*/
@Slf4j
@DubboService(version = "1.0.0", group = "demo", timeout = 3000)
public class UserServiceImpl implements UserService {
private static final List<UserDTO> USER_DATA = Lists.newArrayList();
static {
UserDTO user1 = new UserDTO();
user1.setUserId(1L);
user1.setUsername("jam");
user1.setNickname("果酱");
user1.setEmail("jam@demo.com");
user1.setPhone("13800138000");
user1.setCreateTime(LocalDateTime.now());
user1.setUpdateTime(LocalDateTime.now());
USER_DATA.add(user1);
UserDTO user2 = new UserDTO();
user2.setUserId(2L);
user2.setUsername("ken");
user2.setNickname("Ken");
user2.setEmail("ken@demo.com");
user2.setPhone("13900139000");
user2.setCreateTime(LocalDateTime.now());
user2.setUpdateTime(LocalDateTime.now());
USER_DATA.add(user2);
}
@Override
public Result<UserDTO> getUserById(Long userId) {
log.info("查询用户信息,userId:{}", userId);
if (userId == null || userId <= 0) {
return Result.fail("用户ID不能为空");
}
UserDTO userDTO = USER_DATA.stream()
.filter(user -> userId.equals(user.getUserId()))
.findFirst()
.orElse(null);
return Result.success(userDTO);
}
@Override
public Result<List<UserDTO>> listUserByCondition(UserQueryDTO queryDTO) {
log.info("条件查询用户列表,queryDTO:{}", queryDTO);
if (queryDTO == null) {
return Result.success(USER_DATA);
}
List<UserDTO> resultList = USER_DATA.stream().filter(user -> {
boolean match = true;
if (StringUtils.hasText(queryDTO.getUsername())) {
match = user.getUsername().contains(queryDTO.getUsername());
}
if (match && StringUtils.hasText(queryDTO.getNickname())) {
match = user.getNickname().contains(queryDTO.getNickname());
}
if (match && StringUtils.hasText(queryDTO.getPhone())) {
match = user.getPhone().contains(queryDTO.getPhone());
}
return match;
}).collect(Collectors.toList());
return Result.success(resultList);
}
@Override
public Result<Long> addUser(UserDTO userDTO) {
log.info("新增用户信息,userDTO:{}", userDTO);
if (userDTO == null) {
return Result.fail("用户信息不能为空");
}
if (!StringUtils.hasText(userDTO.getUsername())) {
return Result.fail("用户名不能为空");
}
boolean exist = USER_DATA.stream()
.anyMatch(user -> user.getUsername().equals(userDTO.getUsername()));
if (exist) {
return Result.fail("用户名已存在");
}
Long maxUserId = USER_DATA.stream()
.map(UserDTO::getUserId)
.max(Long::compareTo)
.orElse(0L);
userDTO.setUserId(maxUserId + 1);
userDTO.setCreateTime(LocalDateTime.now());
userDTO.setUpdateTime(LocalDateTime.now());
USER_DATA.add(userDTO);
return Result.success(userDTO.getUserId());
}
}
spring:
application:
name: dubbo-demo-provider
dubbo:
application:
name: ${spring.application.name}
registry-mode: instance
registry:
address: nacos://127.0.0.1:8848
group: demo
protocol:
name: dubbo
port: 20880
serialization: kryo
dispatcher: message
threadpool: cached
threads: 200
iothreads: 8
payload: 8388608
parameters:
tcp.nodelay: true
so.backlog: 1024
so.keepalive: true
provider:
version: 1.0.0
group: demo
timeout: 3000
retries: 0
package com.jam.demo.controller;
import com.jam.demo.dto.UserDTO;
import com.jam.demo.dto.UserQueryDTO;
import com.jam.demo.common.Result;
import com.jam.demo.service.UserService;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.Parameter;
import io.swagger.v3.oas.annotations.tags.Tag;
import lombok.extern.slf4j.Slf4j;
import org.apache.dubbo.config.annotation.DubboReference;
import org.springframework.web.bind.annotation.*;
import java.util.List;
/**
* 用户服务前端控制器
* @author ken
*/
@Slf4j
@RestController
@RequestMapping("/user")
@Tag(name = "用户管理", description = "用户信息管理接口")
public class UserController {
@DubboReference(version = "1.0.0", group = "demo", check = false, timeout = 3000, retries = 2)
private UserService userService;
@GetMapping("/{userId}")
@Operation(summary = "根据用户ID查询用户信息", description = "通过用户ID获取用户详细信息")
public Result<UserDTO> getUserById(
@Parameter(description = "用户ID", required = true, example = "1")
@PathVariable Long userId) {
return userService.getUserById(userId);
}
@PostMapping("/list")
@Operation(summary = "条件查询用户列表", description = "根据查询条件获取用户列表")
public Result<List<UserDTO>> listUserByCondition(@RequestBody UserQueryDTO queryDTO) {
return userService.listUserByCondition(queryDTO);
}
@PostMapping("/add")
@Operation(summary = "新增用户", description = "新增用户信息")
public Result<Long> addUser(@RequestBody UserDTO userDTO) {
return userService.addUser(userDTO);
}
}
spring:
application:
name: dubbo-demo-consumer
server:
port: 8080
dubbo:
application:
name: ${spring.application.name}
registry-mode: instance
registry:
address: nacos://127.0.0.1:8848
group: demo
consumer:
version: 1.0.0
group: demo
timeout: 3000
retries: 2
check: false
springdoc:
api-docs:
enabled: true
swagger-ui:
enabled: true
path: /swagger-ui.html
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.2.4</version>
<relativePath/>
</parent>
<groupId>com.jam.demo</groupId>
<artifactId>grpc-demo</artifactId>
<version>1.0.0</version>
<name>grpc-demo</name>
<properties>
<java.version>17</java.version>
<grpc.version>1.65.1</grpc.version>
<protobuf.version>3.25.3</protobuf.version>
<lombok.version>1.18.30</lombok.version>
<guava.version>33.1.0-jre</guava.version>
</properties>
<dependencies>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-netty-shaded</artifactId>
<version>${grpc.version}</version>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-protobuf</artifactId>
<version>${grpc.version}</version>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-stub</artifactId>
<version>${grpc.version}</version>
</dependency>
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>${protobuf.version}</version>
</dependency>
<dependency>
<groupId>javax.annotation</groupId>
<artifactId>javax.annotation-api</artifactId>
<version>1.3.2</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>${lombok.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
<version>${guava.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<extensions>
<extension>
<groupId>kr.motd.maven</groupId>
<artifactId>os-maven-plugin</artifactId>
<version>1.7.1</version>
</extension>
</extensions>
<plugins>
<plugin>
<groupId>org.xolstice.maven.plugins</groupId>
<artifactId>protobuf-maven-plugin</artifactId>
<version>0.6.1</version>
<configuration>
<protocArtifact>com.google.protobuf:protoc:${protobuf.version}:exe:${os.detected.classifier}</protocArtifact>
<pluginId>grpc-java</pluginId>
<pluginArtifact>io.grpc:protoc-gen-grpc-java:${grpc.version}:exe:${os.detected.classifier}</pluginArtifact>
<protoSourceRoot>src/main/proto</protoSourceRoot>
<outputDirectory>src/main/java</outputDirectory>
<clearOutputDirectory>false</clearOutputDirectory>
</configuration>
<executions>
<execution>
<goals>
<goal>compile</goal>
<goal>compile-custom</goal>
</goals>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<excludes>
<exclude>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</exclude>
</excludes>
</configuration>
</plugin>
</plugins>
</build>
</project>
syntax = "proto3";
option java_multiple_files = true;
option java_package = "com.jam.demo.grpc";
option java_outer_classname = "UserServiceProto";
option objc_class_prefix = "USR";
package user;
import "google/protobuf/timestamp.proto";
// 用户信息
message User {
int64 user_id = 1;
string username = 2;
string nickname = 3;
string email = 4;
string phone = 5;
google.protobuf.Timestamp create_time = 6;
google.protobuf.Timestamp update_time = 7;
}
// 用户查询请求
message UserQueryRequest {
int64 user_id = 1;
}
// 用户列表查询请求
message UserListQueryRequest {
string username = 1;
string nickname = 2;
string phone = 3;
int32 page_num = 4;
int32 page_size = 5;
}
// 用户新增请求
message UserAddRequest {
string username = 1;
string nickname = 2;
string email = 3;
string phone = 4;
}
// 通用响应
message CommonResponse {
int32 code = 1;
string message = 2;
User data = 3;
}
// 用户列表响应
message UserListResponse {
int32 code = 1;
string message = 2;
repeated User data = 3;
}
// 用户新增响应
message UserAddResponse {
int32 code = 1;
string message = 2;
int64 user_id = 3;
}
// 用户服务定义
service UserService {
// 根据用户ID查询用户信息
rpc GetUserById(UserQueryRequest) returns (CommonResponse);
// 条件查询用户列表
rpc ListUserByCondition(UserListQueryRequest) returns (UserListResponse);
// 新增用户信息
rpc AddUser(UserAddRequest) returns (UserAddResponse);
}
package com.jam.demo.service.impl;
import com.google.protobuf.Timestamp;
import com.jam.demo.grpc.*;
import io.grpc.stub.StreamObserver;
import lombok.extern.slf4j.Slf4j;
import org.springframework.util.StringUtils;
import com.google.common.collect.Lists;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.util.List;
import java.util.stream.Collectors;
/**
* gRPC用户服务实现类
* @author ken
*/
@Slf4j
public class UserGrpcServiceImpl extends UserServiceGrpc.UserServiceImplBase {
private static final List<User> USER_DATA = Lists.newArrayList();
static {
Timestamp now = Timestamp.newBuilder()
.setSeconds(Instant.now().getEpochSecond())
.setNanos(Instant.now().getNano())
.build();
User user1 = User.newBuilder()
.setUserId(1L)
.setUsername("jam")
.setNickname("果酱")
.setEmail("jam@demo.com")
.setPhone("13800138000")
.setCreateTime(now)
.setUpdateTime(now)
.build();
USER_DATA.add(user1);
User user2 = User.newBuilder()
.setUserId(2L)
.setUsername("ken")
.setNickname("Ken")
.setEmail("ken@demo.com")
.setPhone("13900139000")
.setCreateTime(now)
.setUpdateTime(now)
.build();
USER_DATA.add(user2);
}
@Override
public void getUserById(UserQueryRequest request, StreamObserver<CommonResponse> responseObserver) {
log.info("gRPC查询用户信息,userId:{}", request.getUserId());
CommonResponse.Builder responseBuilder = CommonResponse.newBuilder();
try {
long userId = request.getUserId();
if (userId <= 0) {
responseBuilder.setCode(500).setMessage("用户ID不能为空");
responseObserver.onNext(responseBuilder.build());
responseObserver.onCompleted();
return;
}
User user = USER_DATA.stream()
.filter(u -> userId == u.getUserId())
.findFirst()
.orElse(null);
responseBuilder.setCode(200).setMessage("操作成功");
if (user != null) {
responseBuilder.setData(user);
}
} catch (Exception e) {
log.error("查询用户信息异常", e);
responseBuilder.setCode(500).setMessage("系统异常");
}
responseObserver.onNext(responseBuilder.build());
responseObserver.onCompleted();
}
@Override
public void listUserByCondition(UserListQueryRequest request, StreamObserver<UserListResponse> responseObserver) {
log.info("gRPC条件查询用户列表,request:{}", request);
UserListResponse.Builder responseBuilder = UserListResponse.newBuilder();
try {
List<User> resultList = USER_DATA.stream().filter(user -> {
boolean match = true;
if (StringUtils.hasText(request.getUsername())) {
match = user.getUsername().contains(request.getUsername());
}
if (match && StringUtils.hasText(request.getNickname())) {
match = user.getNickname().contains(request.getNickname());
}
if (match && StringUtils.hasText(request.getPhone())) {
match = user.getPhone().contains(request.getPhone());
}
return match;
}).collect(Collectors.toList());
responseBuilder.setCode(200).setMessage("操作成功").addAllData(resultList);
} catch (Exception e) {
log.error("查询用户列表异常", e);
responseBuilder.setCode(500).setMessage("系统异常");
}
responseObserver.onNext(responseBuilder.build());
responseObserver.onCompleted();
}
@Override
public void addUser(UserAddRequest request, StreamObserver<UserAddResponse> responseObserver) {
log.info("gRPC新增用户信息,request:{}", request);
UserAddResponse.Builder responseBuilder = UserAddResponse.newBuilder();
try {
if (!StringUtils.hasText(request.getUsername())) {
responseBuilder.setCode(500).setMessage("用户名不能为空");
responseObserver.onNext(responseBuilder.build());
responseObserver.onCompleted();
return;
}
boolean exist = USER_DATA.stream()
.anyMatch(user -> user.getUsername().equals(request.getUsername()));
if (exist) {
responseBuilder.setCode(500).setMessage("用户名已存在");
responseObserver.onNext(responseBuilder.build());
responseObserver.onCompleted();
return;
}
long maxUserId = USER_DATA.stream()
.map(User::getUserId)
.max(Long::compareTo)
.orElse(0L);
long newUserId = maxUserId + 1;
Timestamp now = Timestamp.newBuilder()
.setSeconds(Instant.now().getEpochSecond())
.setNanos(Instant.now().getNano())
.build();
User newUser = User.newBuilder()
.setUserId(newUserId)
.setUsername(request.getUsername())
.setNickname(request.getNickname())
.setEmail(request.getEmail())
.setPhone(request.getPhone())
.setCreateTime(now)
.setUpdateTime(now)
.build();
USER_DATA.add(newUser);
responseBuilder.setCode(200).setMessage("操作成功").setUserId(newUserId);
} catch (Exception e) {
log.error("新增用户信息异常", e);
responseBuilder.setCode(500).setMessage("系统异常");
}
responseObserver.onNext(responseBuilder.build());
responseObserver.onCompleted();
}
}
package com.jam.demo;
import com.jam.demo.service.impl.UserGrpcServiceImpl;
import io.grpc.Server;
import io.grpc.ServerBuilder;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
/**
* gRPC服务端启动类
* @author ken
*/
@Slf4j
@SpringBootApplication
public class GrpcServerApplication {
private Server grpcServer;
private static final int GRPC_PORT = 9090;
public static void main(String[] args) {
SpringApplication.run(GrpcServerApplication.class, args);
}
@PostConstruct
public void startGrpcServer() throws Exception {
ExecutorService businessExecutor = new ThreadPoolExecutor(
40,
200,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(1000),
new ThreadFactory() {
private final AtomicInteger threadNumber = new AtomicInteger(1);
@Override
public Thread newThread(Runnable r) {
return new Thread(r, "grpc-business-thread-" + threadNumber.getAndIncrement());
}
},
new ThreadPoolExecutor.AbortPolicy()
);
grpcServer = ServerBuilder.forPort(GRPC_PORT)
.executor(businessExecutor)
.maxInboundMessageSize(16 * 1024 * 1024)
.maxConcurrentCallsPerConnection(1000)
.addService(new UserGrpcServiceImpl())
.build()
.start();
log.info("gRPC server started on port {}", GRPC_PORT);
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
log.info("shutting down gRPC server");
stopGrpcServer();
}));
}
@PreDestroy
public void stopGrpcServer() {
if (grpcServer != null) {
grpcServer.shutdown();
}
}
}
package com.jam.demo.controller;
import com.jam.demo.grpc.*;
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.Parameter;
import io.swagger.v3.oas.annotations.tags.Tag;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.annotation.*;
import java.util.List;
import java.util.concurrent.TimeUnit;
/**
* gRPC用户服务前端控制器
* @author ken
*/
@Slf4j
@RestController
@RequestMapping("/grpc/user")
@Tag(name = "gRPC用户管理", description = "gRPC用户信息管理接口")
public class GrpcUserController {
private final ManagedChannel channel;
private final UserServiceGrpc.UserServiceBlockingStub userServiceStub;
public GrpcUserController() {
this.channel = ManagedChannelBuilder.forAddress("127.0.0.1", 9090)
.usePlaintext()
.defaultLoadBalancingPolicy("round_robin")
.flowControlWindow(1024 * 1024)
.keepAliveTime(30, TimeUnit.SECONDS)
.keepAliveTimeout(5, TimeUnit.SECONDS)
.keepAliveWithoutCalls(true)
.build();
this.userServiceStub = UserServiceGrpc.newBlockingStub(channel);
}
@GetMapping("/{userId}")
@Operation(summary = "根据用户ID查询用户信息", description = "通过gRPC调用获取用户详细信息")
public CommonResponse getUserById(
@Parameter(description = "用户ID", required = true, example = "1")
@PathVariable Long userId) {
UserQueryRequest request = UserQueryRequest.newBuilder().setUserId(userId).build();
return userServiceStub.getUserById(request);
}
@PostMapping("/list")
@Operation(summary = "条件查询用户列表", description = "通过gRPC调用获取用户列表")
public UserListResponse listUserByCondition(@RequestBody UserListQueryRequest request) {
return userServiceStub.listUserByCondition(request);
}
@PostMapping("/add")
@Operation(summary = "新增用户", description = "通过gRPC调用新增用户信息")
public UserAddResponse addUser(@RequestBody UserAddRequest request) {
return userServiceStub.addUser(request);
}
}
RPC框架是分布式微服务架构的核心基础设施,Dubbo与gRPC作为当前最主流的两款RPC框架,各有其核心优势与适用场景。
Dubbo 3.x 深度适配Java生态,提供了企业级全链路的服务治理能力,在Java为主的微服务架构中,有着天然的优势,适合需要强服务治理、复杂业务场景的企业级应用。
gRPC基于HTTP/2与Protobuf设计,跨语言能力极强,原生支持流式调用,深度适配云原生架构,适合多语言微服务、跨平台服务调用、实时流式数据传输的场景。
在实际选型中,需根据业务的技术栈、场景需求、团队能力综合判断,同时掌握框架的底层原理与调优策略,才能充分发挥RPC框架的性能,构建稳定、高性能的分布式微服务系统。
