在微服务架构中,服务间通信的性能和类型安全直接影响系统整体吞吐和可维护性。gRPC基于HTTP/2和Protocol Buffers,相比REST+JSON在序列化性能、传输效率和接口契约方面优势明显。本文记录Spring Boot微服务架构中集成gRPC的完整方案,涵盖Protobuf定义、服务实现、拦截器治理和生产部署配置。
gRPC vs REST:微服务内部通信选型
微服务内部通信常见的方案有REST/HTTP+JSON、gRPC和消息队列。REST适合面向外部客户端的API,但内部服务间调用在高频场景下JSON序列化开销显著。gRPC的核心优势:
- Protobuf二进制序列化:序列化体积比JSON小3-10倍,反序列化速度快5-100倍
- HTTP/2多路复用:单个TCP连接并发多个请求,消除队头阻塞
- 强类型契约:Protobuf IDL定义接口,自动生成多语言客户端代码
- 内置流式传输:支持服务端流、客户端流和双向流
# 性能对比基准(100万次序列化测试)
# Protobuf: 耗时 0.8s, 体积 28 bytes
# JSON: 耗时 4.3s, 体积 125 bytes
# Protobuf序列化速度约为JSON的5倍,体积约为1/4
Protobuf IDL定义与代码生成
所有gRPC接口先通过Protobuf IDL文件定义。以电商订单服务为例,定义订单查询、创建和流式通知接口。
// proto/order_service.proto
syntax = "proto3";
package com.yunthe.grpc;
option java_multiple_files = true;
option java_package = "com.yunthe.grpc.order";
option java_outer_classname = "OrderServiceProto";
// 订单服务定义
service OrderService {
// 查询订单(一元调用)
rpc GetOrder(GetOrderRequest) returns (OrderResponse);
// 创建订单(一元调用)
rpc CreateOrder(CreateOrderRequest) returns (OrderResponse);
// 批量查询订单(服务端流式)
rpc StreamOrders(StreamOrdersRequest) returns (stream OrderResponse);
// 订单状态变更通知(双向流式)
rpc OrderStatusStream(stream OrderStatusUpdate) returns (stream OrderStatusNotification);
}
message GetOrderRequest {
string order_id = 1;
}
message CreateOrderRequest {
string user_id = 1;
repeated OrderItem items = 2;
string shipping_address = 3;
PaymentMethod payment_method = 4;
}
message OrderItem {
string product_id = 1;
int32 quantity = 2;
double unit_price = 3;
}
message OrderResponse {
string order_id = 1;
string user_id = 2;
OrderStatus status = 3;
repeated OrderItem items = 4;
double total_amount = 5;
int64 created_at = 6; // Unix时间戳
}
message StreamOrdersRequest {
string user_id = 1;
OrderStatus status_filter = 2;
int32 limit = 3;
}
message OrderStatusUpdate {
string order_id = 1;
OrderStatus new_status = 2;
}
message OrderStatusNotification {
string order_id = 1;
OrderStatus previous_status = 2;
OrderStatus current_status = 3;
int64 timestamp = 4;
}
enum OrderStatus {
ORDER_STATUS_UNSPECIFIED = 0;
PENDING = 1;
PAID = 2;
SHIPPED = 3;
DELIVERED = 4;
CANCELLED = 5;
}
enum PaymentMethod {
PAYMENT_METHOD_UNSPECIFIED = 0;
ALIPAY = 1;
WECHAT_PAY = 2;
CREDIT_CARD = 3;
}
Maven项目pom.xml配置protobuf-maven-plugin自动生成Java代码:
<!-- pom.xml -->
<dependencies>
<dependency>
<groupId>net.devh</groupId>
<artifactId>grpc-spring-boot-starter</artifactId>
<version>3.1.0.RELEASE</version>
</dependency>
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>3.25.5</version>
</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:3.25.5:exe:${os.detected.classifier}
</protocArtifact>
<pluginId>grpc-java</pluginId>
<pluginArtifact>
io.grpc:protoc-gen-grpc-java:1.66.0:exe:${os.detected.classifier}
</pluginArtifact>
</configuration>
<executions>
<execution>
<goals>
<goal>compile</goal>
<goal>compile-custom</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
执行mvn compile后,生成的Java类在target/generated-sources/protobuf/目录。
服务端实现与gRPC拦截器
使用grpc-spring-boot-starter后,服务实现只需在Bean上添加@GrpcService注解。拦截器用于统一处理日志、认证、限流等横切关注点。
// 订单服务实现
@GrpcService
public class OrderGrpcService extends OrderServiceGrpc.OrderServiceImplBase {
private final OrderRepository orderRepository;
private final OrderEventPublisher eventPublisher;
public OrderGrpcService(OrderRepository orderRepository,
OrderEventPublisher eventPublisher) {
this.orderRepository = orderRepository;
this.eventPublisher = eventPublisher;
}
@Override
public void getOrder(GetOrderRequest request,
StreamObserver<OrderResponse> responseObserver) {
try {
Order order = orderRepository.findById(request.getOrderId())
.orElseThrow(() -> Status.NOT_FOUND
.withDescription("订单不存在: " + request.getOrderId())
.asRuntimeException());
OrderResponse response = toProtoResponse(order);
responseObserver.onNext(response);
responseObserver.onCompleted();
} catch (Exception e) {
responseObserver.onError(e);
}
}
@Override
public void createOrder(CreateOrderRequest request,
StreamObserver<OrderResponse> responseObserver) {
// 参数校验
if (request.getUserId().isEmpty()) {
responseObserver.onError(Status.INVALID_ARGUMENT
.withDescription("user_id不能为空")
.asRuntimeException());
return;
}
if (request.getItemsCount() == 0) {
responseObserver.onError(Status.INVALID_ARGUMENT
.withDescription("订单必须包含至少一个商品")
.asRuntimeException());
return;
}
// 创建订单
Order order = Order.builder()
.userId(request.getUserId())
.items(request.getItemsList().stream()
.map(this::toOrderItem)
.collect(Collectors.toList()))
.shippingAddress(request.getShippingAddress())
.status(OrderStatus.PENDING)
.build();
// 计算总金额
order.calculateTotal();
Order saved = orderRepository.save(order);
eventPublisher.publishOrderCreated(saved);
responseObserver.onNext(toProtoResponse(saved));
responseObserver.onCompleted();
}
@Override
public void streamOrders(StreamOrdersRequest request,
StreamObserver<OrderResponse> responseObserver) {
// 服务端流式:分批发送订单
try (Stream<Order> orderStream = orderRepository
.streamByUserIdAndStatus(
request.getUserId(),
toEntityStatus(request.getStatusFilter()),
PageRequest.of(0, request.getLimit() > 0 ? request.getLimit() : 100)
)) {
orderStream.forEach(order -> {
responseObserver.onNext(toProtoResponse(order));
});
responseObserver.onCompleted();
} catch (Exception e) {
responseObserver.onError(Status.INTERNAL
.withDescription("流式查询失败: " + e.getMessage())
.asRuntimeException());
}
}
private OrderResponse toProtoResponse(Order order) {
return OrderResponse.newBuilder()
.setOrderId(order.getId())
.setUserId(order.getUserId())
.setStatus(toProtoStatus(order.getStatus()))
.addAllItems(order.getItems().stream()
.map(item -> OrderItem.newBuilder()
.setProductId(item.getProductId())
.setQuantity(item.getQuantity())
.setUnitPrice(item.getUnitPrice())
.build())
.collect(Collectors.toList()))
.setTotalAmount(order.getTotalAmount())
.setCreatedAt(order.getCreatedAt().getEpochSecond())
.build();
}
}
全局拦截器实现认证和日志记录:
// 服务端拦截器:认证 + 日志
@GrpcGlobalServerInterceptor
public class AuthAndLogInterceptor implements ServerInterceptor {
private static final Logger log = LoggerFactory.getLogger(AuthAndLogInterceptor.class);
private static final Metadata.Key<String> AUTH_TOKEN =
Metadata.Key.of("authorization", Metadata.ASCII_STRING_MARSHALLER);
private final TokenValidator tokenValidator;
public AuthAndLogInterceptor(TokenValidator tokenValidator) {
this.tokenValidator = tokenValidator;
}
@Override
public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(
ServerCall<ReqT, RespT> call,
Metadata headers,
ServerCallHandler<ReqT, RespT> next) {
String methodName = call.getMethodDescriptor().getFullMethodName();
long startTime = System.nanoTime();
// 内部健康检查跳过认证
if (methodName.contains("HealthCheck")) {
return next.startCall(call, headers);
}
// 认证校验
String token = headers.get(AUTH_TOKEN);
if (token == null || !tokenValidator.validate(token)) {
call.close(Status.UNAUTHENTICATED
.withDescription("无效或缺失的认证token"), headers);
return new ServerCall.Listener<>() {};
}
// 记录请求日志
log.info("gRPC请求: method={}, peer={}",
methodName, call.getAttributes().get(Grpc.TRANSPORT_ATTR_REMOTE_ADDR));
// 包装Listener用于记录响应
ServerCall.Listener<ReqT> delegate = next.startCall(call, headers);
return new ForwardingServerCallListener.SimpleForwardingServerCallListener<>(delegate) {
@Override
public void onComplete() {
long duration = (System.nanoTime() - startTime) / 1_000_000;
log.info("gRPC完成: method={}, 耗时={}ms", methodName, duration);
super.onComplete();
}
@Override
public void onCancel() {
long duration = (System.nanoTime() - startTime) / 1_000_000;
log.warn("gRPC取消: method={}, 耗时={}ms", methodName, duration);
super.onCancel();
}
};
}
}
客户端调用与负载均衡
gRPC客户端使用@GrpcClient注入Stub,grpc-spring-boot-starter自动管理Channel生命周期。服务发现层面可集成Nacos/Consul实现客户端负载均衡。
// gRPC客户端配置
@Configuration
public class GrpcClientConfig {
@Bean
@GrpcClient("order-service")
public OrderServiceGrpc.OrderServiceBlockingStub orderBlockingStub() {
return OrderServiceGrpc.newBlockingStub(
ManagedChannelBuilder.forAddress("nacos:///order-service")
.defaultLoadBalancingPolicy("round_robin")
.intercept(new ClientLogInterceptor())
.build()
);
}
}
// 服务消费者调用
@Service
public class OrderQueryService {
private final OrderServiceGrpc.OrderServiceBlockingStub orderStub;
public OrderQueryService(
@GrpcClient("order-service") OrderServiceGrpc.OrderServiceBlockingStub orderStub) {
this.orderStub = orderStub;
}
// 一元调用
public OrderDTO getOrder(String orderId) {
GetOrderRequest request = GetOrderRequest.newBuilder()
.setOrderId(orderId)
.build();
// 设置超时
OrderResponse response = orderStub
.withDeadlineAfter(3, TimeUnit.SECONDS)
.getOrder(request);
return toDTO(response);
}
// 服务端流式调用
public List<OrderDTO> batchQueryOrders(String userId, int limit) {
StreamOrdersRequest request = StreamOrdersRequest.newBuilder()
.setUserId(userId)
.setLimit(limit)
.build();
List<OrderDTO> results = new ArrayList<>();
Iterator<OrderResponse> iterator = orderStub
.withDeadlineAfter(30, TimeUnit.SECONDS)
.streamOrders(request);
while (iterator.hasNext()) {
results.add(toDTO(iterator.next()));
}
return results;
}
// 双向流式调用
public void subscribeOrderStatus(List<String> orderIds) {
OrderServiceGrpc.OrderServiceStub asyncStub = OrderServiceGrpc.newStub(
orderStub.getChannel()
);
StreamObserver<OrderStatusUpdate> requestObserver = asyncStub
.orderStatusStream(new StreamObserver<>() {
@Override
public void onNext(OrderStatusNotification notification) {
log.info("订单状态变更: {} {} -> {}",
notification.getOrderId(),
notification.getPreviousStatus(),
notification.getCurrentStatus());
}
@Override
public void onError(Throwable t) {
log.error("订单状态流异常", t);
}
@Override
public void onCompleted() {
log.info("订单状态流结束");
}
});
// 发送状态更新请求
orderIds.forEach(id -> requestObserver.onNext(
OrderStatusUpdate.newBuilder()
.setOrderId(id)
.setNewStatus(OrderStatus.SHIPPED)
.build()
));
requestObserver.onCompleted();
}
}
生产环境配置与健康检查
Spring Boot application.yml中的gRPC配置,包含端口、线程池、消息大小限制和健康检查:
# application.yml
grpc:
server:
port: 9090
security:
enabled: true
certificateChain: classpath:certs/server.crt
privateKey: classpath:certs/server.key
max-inbound-message-size: 16MB
max-inbound-metadata-size: 8KB
reflection-service-enabled: true # 开发环境开启反射
client:
order-service:
address: 'nacos:///order-service'
negotiation-type: plaintext # 生产环境改为TLS
max-retry-attempts: 3
default-load-balancing-policy: round_robin
# gRPC健康检查实现
@Service
public class GrpcHealthService extends HealthGrpc.HealthImplBase {
private final HealthIndicator[] healthIndicators;
public GrpcHealthService(HealthIndicator[] healthIndicators) {
this.healthIndicators = healthIndicators;
}
@Override
public void check(HealthCheckRequest request,
StreamObserver<HealthCheckResponse> responseObserver) {
HealthCheckResponse.ServingStatus status = ServingStatus.SERVING;
for (HealthIndicator indicator : healthIndicators) {
Health health = indicator.health();
if (health.getStatus() != Status.UP) {
status = ServingStatus.NOT_SERVING;
break;
}
}
responseObserver.onNext(HealthCheckResponse.newBuilder()
.setStatus(status)
.build());
responseObserver.onCompleted();
}
}
gRPC在微服务内部通信场景中相比REST有明显性能优势,Protobuf的强类型契约也减少了接口对接成本。落地时需注意几点:Protobuf字段编号一旦发布不可变更(向后兼容性)、生产环境务必启用TLS、超时和重试策略需按业务场景精细化配置、gRPC反射服务在生产环境关闭以减少攻击面。服务治理层面可结合Nacos实现服务注册发现与动态负载均衡,配合SkyWalking做gRPC链路追踪,构建完整的微服务可观测性体系。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-wei-fu-wu-grpc-tong-xin-shi-zhan-protobuf-ding/