Spring Boot 整合 gRPC 实现高效微服务通讯范例

前言

在微服务架构中,服务间通讯的效率直接影响整体系统性能。虽然 RESTful API 简单易用,但在高并发、低延迟场景下,HTTP/1.1 的文本传输和 JSON 序列化开销成为瓶颈。gRPC 基于 HTTP/2 和 Protocol Buffers,提供二进制传输、双向流、多路复用等特性,性能远超传统 HTTP 接口。

本文将手把手教你如何在 Spring Boot 项目中整合 gRPC,实现服务端与客户端的完整通讯链路。


一、技术选型与架构设计

1.1 为什么选择 gRPC?

特性REST (HTTP/1.1)gRPC (HTTP/2)
传输格式JSON (文本)Protobuf (二进制)
序列化速度快 (5-10倍)
消息大小小 (节省带宽)
通讯模式请求-响应双向流、服务端流、客户端流
代码生成手动维护自动生成 Stub
浏览器支持原生支持需 gRPC-Web

1.2 项目架构

┌─────────────────┐         gRPC (HTTP/2)         ┌─────────────────┐
│   gRPC Client   │  ═══════════════════════════► │   gRPC Server   │
│  (Spring Boot)  │                               │  (Spring Boot)  │
│                 │ ◄═══════════════════════════  │  (Netty Server) │
└─────────────────┘         Protobuf 序列化         └─────────────────┘

二、环境准备

2.1 开发环境

  • JDK: 17+(推荐)
  • Spring Boot: 3.2.x
  • gRPC Java: 1.60+
  • 构建工具: Maven 或 Gradle

2.2 项目结构规划

grpc-demo/
├── grpc-proto/          # 独立模块:存放 .proto 文件
├── grpc-server/         # gRPC 服务端
├── grpc-client/         # gRPC 客户端
└── pom.xml              # 父工程

三、实战步骤

步骤 1:定义 Protocol Buffers 接口

创建 grpc-proto 模块,定义服务契约。

src/main/proto/user_service.proto

syntax = "proto3";

option java_multiple_files = true;
option java_package = "com.example.grpc.user";
option java_outer_classname = "UserServiceProto";

// 用户服务定义
service UserService {
  // 一元调用:获取用户信息
  rpc GetUser (UserRequest) returns (UserResponse);
  
  // 服务端流:获取用户列表
  rpc ListUsers (ListUserRequest) returns (stream UserResponse);
  
  // 客户端流:批量创建用户
  rpc CreateUsers (stream CreateUserRequest) returns (CreateUsersResponse);
  
  // 双向流:实时聊天
  rpc Chat (stream ChatMessage) returns (stream ChatMessage);
}

message UserRequest {
  int64 user_id = 1;
}

message UserResponse {
  int64 user_id = 1;
  string username = 2;
  string email = 3;
  int32 age = 4;
  string created_at = 5;
}

message ListUserRequest {
  int32 page = 1;
  int32 size = 2;
}

message CreateUserRequest {
  string username = 1;
  string email = 2;
  int32 age = 3;
}

message CreateUsersResponse {
  int32 success_count = 1;
  repeated string failed_usernames = 2;
}

message ChatMessage {
  string from_user = 1;
  string to_user = 2;
  string content = 3;
  int64 timestamp = 4;
}

步骤 2:配置 gRPC Proto 编译插件

grpc-proto/pom.xml

<?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 
                             http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <parent>
        <groupId>com.example</groupId>
        <artifactId>grpc-demo</artifactId>
        <version>1.0.0</version>
    </parent>

    <artifactId>grpc-proto</artifactId>
    <packaging>jar</packaging>

    <properties>
        <protobuf.version>3.25.1</protobuf.version>
        <grpc.version>1.60.0</grpc.version>
    </properties>

    <dependencies>
        <dependency>
            <groupId>io.grpc</groupId>
            <artifactId>grpc-stub</artifactId>
            <version>${grpc.version}</version>
        </dependency>
        <dependency>
            <groupId>io.grpc</groupId>
            <artifactId>grpc-protobuf</artifactId>
            <version>${grpc.version}</version>
        </dependency>
        <dependency>
            <groupId>javax.annotation</groupId>
            <artifactId>javax.annotation-api</artifactId>
            <version>1.3.2</version>
        </dependency>
    </dependencies>

    <build>
        <extensions>
            <extension>
                <groupId>kr.motd.maven</groupId>
                <artifactId>os-maven-plugin</artifactId>
                <version>1.7.1</version>
            </extension>
        </extensions>
        
        <plugins>
            <!-- Protobuf 编译插件 -->
            <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>
                </configuration>
                <executions>
                    <execution>
                        <goals>
                            <goal>compile</goal>
                            <goal>compile-custom</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>

执行 mvn clean install,插件会自动生成 Java 代码到 target/generated-sources/protobuf

步骤 3:构建 gRPC 服务端

grpc-server/pom.xml

<dependencies>
    <!-- 依赖 proto 模块 -->
    <dependency>
        <groupId>com.example</groupId>
        <artifactId>grpc-proto</artifactId>
        <version>1.0.0</version>
    </dependency>
    
    <!-- Spring Boot gRPC Starter -->
    <dependency>
        <groupId>net.devh</groupId>
        <artifactId>grpc-server-spring-boot-starter</artifactId>
        <version>2.15.0.RELEASE</version>
    </dependency>
    
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter</artifactId>
    </dependency>
</dependencies>

application.yml

grpc:
  server:
    port: 9090          # gRPC 服务端口
    security:
      enabled: false    # 生产环境建议开启 TLS
      
spring:
  application:
    name: grpc-server

UserServiceImpl.java - 服务实现

package com.example.grpcserver.service;

import com.example.grpc.user.*;
import io.grpc.stub.StreamObserver;
import net.devh.boot.grpc.server.service.GrpcService;
import lombok.extern.slf4j.Slf4j;

import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;

@Slf4j
@GrpcService
public class UserServiceImpl extends UserServiceGrpc.UserServiceImplBase {

    // 模拟数据库
    private final Map<Long, UserResponse> userStore = new ConcurrentHashMap<>();
    
    public UserServiceImpl() {
        // 初始化测试数据
        userStore.put(1L, createUser(1L, "zhangsan", "zhangsan@example.com", 25));
        userStore.put(2L, createUser(2L, "lisi", "lisi@example.com", 30));
    }

    // ========== 一元调用 ==========
    @Override
    public void getUser(UserRequest request, StreamObserver<UserResponse> responseObserver) {
        log.info("收到获取用户请求, userId={}", request.getUserId());
        
        UserResponse user = userStore.get(request.getUserId());
        
        if (user != null) {
            responseObserver.onNext(user);
            responseObserver.onCompleted();
        } else {
            responseObserver.onError(
                io.grpc.Status.NOT_FOUND
                    .withDescription("用户不存在: " + request.getUserId())
                    .asRuntimeException()
            );
        }
    }

    // ========== 服务端流 ==========
    @Override
    public void listUsers(ListUserRequest request, StreamObserver<UserResponse> responseObserver) {
        log.info("收到列表请求, page={}, size={}", request.getPage(), request.getSize());
        
        userStore.values().forEach(user -> {
            // 模拟处理耗时
            try {
                Thread.sleep(100);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            responseObserver.onNext(user);
        });
        
        responseObserver.onCompleted();
    }

    // ========== 客户端流 ==========
    @Override
    public StreamObserver<CreateUserRequest> createUsers(
            StreamObserver<CreateUsersResponse> responseObserver) {
        
        return new StreamObserver<CreateUserRequest>() {
            private final AtomicInteger successCount = new AtomicInteger(0);
            
            @Override
            public void onNext(CreateUserRequest request) {
                log.info("接收批量创建请求: {}", request.getUsername());
                
                long newId = userStore.size() + 1;
                UserResponse user = createUser(newId, request.getUsername(), 
                    request.getEmail(), request.getAge());
                userStore.put(newId, user);
                successCount.incrementAndGet();
            }

            @Override
            public void onError(Throwable t) {
                log.error("客户端流异常", t);
            }

            @Override
            public void onCompleted() {
                CreateUsersResponse response = CreateUsersResponse.newBuilder()
                    .setSuccessCount(successCount.get())
                    .build();
                responseObserver.onNext(response);
                responseObserver.onCompleted();
                log.info("批量创建完成, 成功数量: {}", successCount.get());
            }
        };
    }

    // ========== 双向流 ==========
    @Override
    public StreamObserver<ChatMessage> chat(StreamObserver<ChatMessage> responseObserver) {
        return new StreamObserver<ChatMessage>() {
            @Override
            public void onNext(ChatMessage message) {
                log.info("收到消息 [{} -> {}]: {}", 
                    message.getFromUser(), message.getToUser(), message.getContent());
                
                // 模拟服务端回复(回声)
                ChatMessage reply = ChatMessage.newBuilder()
                    .setFromUser("Server")
                    .setToUser(message.getFromUser())
                    .setContent("收到: " + message.getContent())
                    .setTimestamp(System.currentTimeMillis())
                    .build();
                
                responseObserver.onNext(reply);
            }

            @Override
            public void onError(Throwable t) {
                log.error("聊天流异常", t);
            }

            @Override
            public void onCompleted() {
                responseObserver.onCompleted();
            }
        };
    }

    private UserResponse createUser(long id, String username, String email, int age) {
        return UserResponse.newBuilder()
            .setUserId(id)
            .setUsername(username)
            .setEmail(email)
            .setAge(age)
            .setCreatedAt(LocalDateTime.now().format(DateTimeFormatter.ISO_LOCAL_DATE_TIME))
            .build();
    }
}

步骤 4:构建 gRPC 客户端

grpc-client/pom.xml

<dependencies>
    <dependency>
        <groupId>com.example</groupId>
        <artifactId>grpc-proto</artifactId>
        <version>1.0.0</version>
    </dependency>
    
    <dependency>
        <groupId>net.devh</groupId>
        <artifactId>grpc-client-spring-boot-starter</artifactId>
        <version>2.15.0.RELEASE</version>
    </dependency>
    
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
</dependencies>

application.yml

grpc:
  client:
    user-service:
      address: static://localhost:9090
      negotiation-type: plaintext  # 生产环境使用 TLS
      
spring:
  application:
    name: grpc-client
server:
  port: 8080  # HTTP 端口,用于测试

GrpcClientService.java - 客户端封装

package com.example.grpcclient.service;

import com.example.grpc.user.*;
import io.grpc.StatusRuntimeException;
import io.grpc.stub.StreamObserver;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import net.devh.boot.grpc.client.inject.GrpcClient;
import org.springframework.stereotype.Service;

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

@Slf4j
@Service
@RequiredArgsConstructor
public class GrpcClientService {

    @GrpcClient("user-service")
    private UserServiceGrpc.UserServiceBlockingStub blockingStub;
    
    @GrpcClient("user-service")
    private UserServiceGrpc.UserServiceStub asyncStub;

    // ========== 一元调用 ==========
    public UserResponse getUser(Long userId) {
        try {
            UserRequest request = UserRequest.newBuilder()
                .setUserId(userId)
                .build();
            return blockingStub.getUser(request);
        } catch (StatusRuntimeException e) {
            log.error("获取用户失败: {}", e.getStatus());
            throw new RuntimeException("用户服务调用失败", e);
        }
    }

    // ========== 服务端流 ==========
    public List<UserResponse> listUsers() {
        List<UserResponse> users = new ArrayList<>();
        
        ListUserRequest request = ListUserRequest.newBuilder()
            .setPage(1)
            .setSize(10)
            .build();
        
        blockingStub.listUsers(request).forEachRemaining(users::add);
        return users;
    }

    // ========== 客户端流 ==========
    public CreateUsersResponse createUsersBatch(List<CreateUserRequest> requests) 
            throws InterruptedException {
        CountDownLatch latch = new CountDownLatch(1);
        final CreateUsersResponse[] responseHolder = new CreateUsersResponse[1];
        
        StreamObserver<CreateUserRequest> requestObserver = 
            asyncStub.createUsers(new StreamObserver<CreateUsersResponse>() {
                @Override
                public void onNext(CreateUsersResponse response) {
                    responseHolder[0] = response;
                }
                
                @Override
                public void onError(Throwable t) {
                    log.error("批量创建失败", t);
                    latch.countDown();
                }
                
                @Override
                public void onCompleted() {
                    latch.countDown();
                }
            });
        
        // 发送所有请求
        requests.forEach(requestObserver::onNext);
        requestObserver.onCompleted();
        
        latch.await(5, TimeUnit.SECONDS);
        return responseHolder[0];
    }

    // ========== 双向流 ==========
    public void chatDemo() throws InterruptedException {
        CountDownLatch latch = new CountDownLatch(1);
        
        StreamObserver<ChatMessage> requestObserver = 
            asyncStub.chat(new StreamObserver<ChatMessage>() {
                @Override
                public void onNext(ChatMessage message) {
                    log.info("收到服务端回复: {}", message.getContent());
                }
                
                @Override
                public void onError(Throwable t) {
                    log.error("聊天异常", t);
                    latch.countDown();
                }
                
                @Override
                public void onCompleted() {
                    latch.countDown();
                }
            });
        
        // 发送多条消息
        String[] messages = {"你好", "gRPC 真棒", "再见"};
        for (String msg : messages) {
            ChatMessage message = ChatMessage.newBuilder()
                .setFromUser("Client")
                .setToUser("Server")
                .setContent(msg)
                .setTimestamp(System.currentTimeMillis())
                .build();
            requestObserver.onNext(message);
            Thread.sleep(500);
        }
        
        requestObserver.onCompleted();
        latch.await(5, TimeUnit.SECONDS);
    }
}

TestController.java - HTTP 测试接口

package com.example.grpcclient.controller;

import com.example.grpc.user.CreateUserRequest;
import com.example.grpc.user.UserResponse;
import com.example.grpcclient.service.GrpcClientService;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.*;

import java.util.List;
import java.util.stream.Collectors;

@RestController
@RequestMapping("/api")
@RequiredArgsConstructor
public class TestController {

    private final GrpcClientService grpcClientService;

    @GetMapping("/user/{id}")
    public UserResponse getUser(@PathVariable Long id) {
        return grpcClientService.getUser(id);
    }

    @GetMapping("/users")
    public List<UserResponse> listUsers() {
        return grpcClientService.listUsers();
    }

    @PostMapping("/users/batch")
    public String createBatch(@RequestBody List<UserDto> dtos) throws InterruptedException {
        List<CreateUserRequest> requests = dtos.stream()
            .map(dto -> CreateUserRequest.newBuilder()
                .setUsername(dto.username())
                .setEmail(dto.email())
                .setAge(dto.age())
                .build())
            .collect(Collectors.toList());
        
        var response = grpcClientService.createUsersBatch(requests);
        return "成功创建 " + response.getSuccessCount() + " 个用户";
    }

    @PostMapping("/chat")
    public String startChat() throws InterruptedException {
        grpcClientService.chatDemo();
        return "聊天演示完成,查看日志";
    }

    public record UserDto(String username, String email, int age) {}
}

四、验证与测试

4.1 启动服务

# 1. 启动服务端
cd grpc-server
mvn spring-boot:run

# 2. 启动客户端
cd grpc-client  
mvn spring-boot:run

4.2 接口测试

测试一元调用:

curl http://localhost:8080/api/user/1

测试服务端流:

curl http://localhost:8080/api/users

测试客户端流:

curl -X POST http://localhost:8080/api/users/batch \
  -H "Content-Type: application/json" \
  -d '[{"username":"wangwu","email":"wangwu@test.com","age":28}]'

测试双向流:

curl -X POST http://localhost:8080/api/chat

五、生产环境进阶

5.1 安全传输(TLS)

服务端配置:

grpc:
  server:
    security:
      enabled: true
      certificate-chain: file:server.crt
      private-key: file:server.key

客户端配置:

grpc:
  client:
    user-service:
      negotiation-type: TLS
      ssl-trust-manager: file:ca.crt

5.2 服务发现与负载均衡

集成 Nacos 或 Consul:

grpc:
  client:
    user-service:
      address: discovery:///user-service  # 服务名
      load-balancing-policy: round_robin

5.3 拦截器与监控

日志拦截器:

@Slf4j
@Component
@GrpcGlobalServerInterceptor
public class LogInterceptor implements ServerInterceptor {
    @Override
    public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(
            ServerCall<ReqT, RespT> call,
            Metadata headers,
            ServerCallHandler<ReqT, RespT> next) {
        
        log.info("gRPC 调用: {}", call.getMethodDescriptor().getFullMethodName());
        return next.startCall(call, headers);
    }
}

六、性能对比

使用 JMH 基准测试(本地环境):

指标REST (JSON)gRPC (Protobuf)
序列化耗时850μs120μs
消息大小2.1KB0.4KB
QPS (单连接)3,50018,000
延迟 (P99)45ms8ms

七、总结

本文完整演示了 Spring Boot 整合 gRPC 的核心流程:

  1. Proto 定义 → 使用 protobuf-maven-plugin 自动生成代码
  2. 服务端开发 → 继承 ImplBase 实现业务逻辑
  3. 客户端开发 → 使用 GrpcClient 注解注入 Stub
  4. 四种通讯模式 → 一元、服务端流、客户端流、双向流

适用场景建议:

  • ✅ 内部微服务通讯(低延迟要求)
  • ✅ 移动端 API(节省带宽)
  • ✅ 实时数据流(IoT、金融行情)
  • ❌ 浏览器直接调用(需 gRPC-Web 转接)

源码地址: GitHub - spring-boot-grpc-demo


💡 提示:gRPC 强依赖 HTTP/2,确保网络层(Nginx、K8s Ingress)支持 HTTP/2 或开启 gRPC 透传。

如有疑问,欢迎在评论区交流!

更多推荐