
最近在搞一个基于 Vert.x 的实时消息推送服务第一次把 gRPC 用在生产环境踩了不少坑。网上讲 Vert.x 和 gRPC 的教程不少但多数偏向 hello world真正落到“消息发送”这个场景从 proto 设计到连接管理、超时重试、背压控制坑其实比预想的多。我把自己实际折腾完的完整方案整理出来从服务端到客户端附带可直接复制的代码和参数配置给正打算用 Vert.x 接 gRPC 做消息推送的朋友做一个参考。这套方案适合两类人一是已经在用 Vert.x 做基础框架想把 RPC 通信从 REST HTTP 切换到 gRPC二是刚开始接触 gRPC想知道在一个响应式框架里该怎么优雅地收发消息、处理流式数据。即便你完全没碰过 gRPC只要跟着后面的步骤走也能把一条消息从客户端发到服务端再拿到确认回执。1. 项目整体设计与思路拆解1.1 消息发送场景的需求拆解一个消息发送系统从底层看无非是几件事建立连接、把消息体和元信息序列化、通过网络传输、在接收端反序列化、给发送方一个确认。但真正设计时会发现“消息”这两个字背后有完全不同的含义。拿我这个项目来说既有单条推送比如用户触发一条改密通知也有批量下发比如运营要群发模板消息还有连续的流式转发比如日志实时上报。这三种模式在 REST 里往往要用三套接口在 gRPC 里其实是一个服务的三种方法类型。所以在动手前我先把场景拆成了三类一元消息请求一次响应一次最典型的发送确认模型。客户端流式客户端连续发送多条消息服务端统一返回汇总结果适合批量上报。双向流式建立一条长连接双方各自独立发送适合聊天、实时推送。这个拆解决定了 proto 文件里的 service 定义也决定了后面客户端调用的 API 形态。很多人一开始就扑上去写代码结果做到一半发现模式选错了回头改 proto 成本很高。这不是小事后续牵扯到代码生成、测试 mock、AB 兼容所以花点时间把模式定死是最值得的一步。1.2 为什么选 gRPC 而不是 REST 或 WebSocket这个选择我很明确。REST 的优点是调试方便、生态成熟但做消息推送时有两个很痛的点一是接口没有强约束参数靠文档和自觉前后端经常为字段对不上扯皮二是不擅长流式SSE 虽然可以做服务端推流但客户端往服务端持续灌数据就很难受。WebSocket 是另一个常用选项它是双向长连接适合聊天类场景但没有一套自动化的、跨语言的 IDL 和序列化方案所有消息格式都得自己定义自己解析传输效率也只能靠自己优化。也有朋友问过干脆用 TCP 长连接裸发不就行了如果只是单个设备上报比如 esp01s 这种资源极其有限的终端裸 TCP 当然够用但到了服务端多服务、多语言互调的阶段裸 TCP 没有 IDL 约束、没有跨语言序列化每个客户端都要自造协议维护成本会直线上升。gRPC 的价值就在这协议的通用性和传输的高效性能同时拿到。gRPC 恰好把 REST 和 WebSocket 缺的部分都补上了。它用 protobuf 做二进制序列化字段有强约束新增字段还能做到向后兼容底层传输是基于 HTTP/2 的多路复用一条 TCP 连接上可以同时跑很多个请求不会像 HTTP/1.1 那样被队头阻塞拖住四种 RPC 模式天然覆盖了从单条发送到双向流的所有场景。对 Vert.x 来说它本身是事件驱动、响应式模型gRPC 基于 HTTP/2 的异步交互和它非常契合不像传统阻塞式 Java 服务里用 gRPC 还需要考虑额外的线程池。当然gRPC 也有学习成本和调试成本protobuf 的二进制报文肉眼不可读抓包要配专门的工具。所以我的结论是对内模块之间、跨语言服务之间、需要高吞吐低延迟的发送链路用 gRPC 值得对外浏览器环境的简单 REST 接口该保留还保留。两手都用不矛盾。1.3 与 Vert.x 事件循环模型的配合方式这一点我觉得是 Vert.x 下用 gRPC 最容易被忽略的地方。Vert.x 的核心是 Event Loop所有 IO 事件都在这个线程上处理任何阻塞操作都会卡住整个事件循环。传统的 gRPC Java 客户端 stub 分 Blocking阻塞式和 Async异步式如果在 Vert.x 的 Event Loop 里直接调用 BlockingStub面向连接操作的线程调度会直接卡你的 Event Loop表现就是接口莫名其妙延迟暴涨、CPU 不高但吞吐量上不去。所以设计原则是在 Event Loop 中全部使用异步 stub用回调或 Future 处理结果。需要做阻塞操作比如 DB 查询、文件导出时用vertx.executeBlocking或 worker verticle 隔离。Channel 和 Stub 的创建是重量级操作尽量复用不要每次发送消息都新建。我自己的做法是启动时创建好GrpcClientChannel和 stub用单例持有发送方法里只做 request 构建和回调注册。后面的代码也是这样处理的先把这个原则立住写代码才不会走偏。2. 环境准备与依赖配置2.1 Maven 依赖与版本对齐Vert.x 的 gRPC 支持在 4.x 开始已经比较成熟版本上建议使用 4.4 以上我用的 4.5.8。这里要特别提醒vertx-grpc自身依赖了 grpc-java 的接口但实现还是由 grpc-netty 提供的所以必须显式引入grpc-netty-shaded或grpc-netty。如果你只加了vertx-grpc运行时会报找不到类。pom.xml 里核心依赖这样写properties vertx.version4.5.8/vertx.version grpc.version1.62.2/grpc.version protobuf.version3.25.3/protobuf.version /properties dependencies dependency groupIdio.vertx/groupId artifactIdvertx-grpc/artifactId version${vertx.version}/version /dependency dependency groupIdio.grpc/groupId artifactIdgrpc-netty-shaded/artifactId version${grpc.version}/version /dependency dependency groupIdio.grpc/groupId artifactIdgrpc-protobuf/artifactId version${grpc.version}/version /dependency dependency groupIdio.grpc/groupId artifactIdgrpc-stub/artifactId version${grpc.version}/version /dependency dependency groupIdjakarta.annotation/groupId artifactIdjakarta.annotation-api/artifactId version1.3.5/version /dependency /dependencies关于grpc-netty-shaded和grpc-netty的选择shaded 版本会把 netty 相关的类重新打包改名避免和你项目里其他 netty 版本冲突。Vert.x 本身内置了 netty版本通常是 4.1.x如果直接引入grpc-netty自带的 netty 依赖很可能出现两个 netty 版本在 classpath 里打架。我建议先上 shaded省心。如果你的项目里已经明确把 netty 版本统一到一个较新版本再用grpc-netty。2.2 proto 文件设计与字段规划发消息的 proto 我设计得比较克制没有一开始就把所有业务字段塞进去。消息发送拆成两个核心部分路由元信息和消息内容。路由元信息包括消息 ID、发送方、接收方、时间戳消息内容就是 content 字段实际业务中可以是 JSON 字符串或者二进制数据。syntax proto3; package msg; option java_multiple_files true; option java_package com.example.grpc.msg; option java_outer_classname MsgProto; service MessageService { rpc SendMessage (SendRequest) returns (SendReply); rpc SendMessageStream (stream SendRequest) returns (stream SendReply); } message SendRequest { string messageId 1; string from 2; string to 3; string content 4; int64 timestamp 5; } message SendReply { bool success 1; string message 2; string messageId 3; string traceId 4; }这里几个字段的设计有个细节timestamp 我特意用 int64 而不是 string因为 proto3 的 JSON 映射里 int64 会默认输出成字符串但二进制传输时是变长整型效率更高。messageId 是必填的后面做发送成功确认和幂等都靠它。traceId 放在 reply 里方便把一次发送链路串起来排查。如果你有多个服务要复用这套消息模型可以把 message 定义单独放到一个 proto 文件service 定义放另一个 proto 文件用 import 引入。我目前项目规模还没到这个程度先把一个文件做干净比拆文件更有价值。2.3 使用 protobuf-maven-plugin 生成代码生成代码我建议用 protobuf-maven-plugin 自动完成不要手动下 protoc 编译否则团队里每个人环境不同生成的代码八成都对不上。build extensions extension groupIdkr.motd.maven/groupId artifactIdos-maven-plugin/artifactId version1.7.1/version /extension /extensions plugins plugin groupIdorg.xolstice.maven.plugins/groupId artifactIdprotobuf-maven-plugin/artifactId version0.6.1/version configuration protocArtifactcom.google.protobuf:protoc:3.25.3:exe:${os.detected.classifier}/protocArtifact pluginIdgrpc-java/pluginId pluginArtifactio.grpc:protoc-gen-grpc-java:1.62.2:exe:${os.detected.classifier}/pluginArtifact /configuration executions execution goals goalcompile/goal goalcompile-custom/goal /goals /execution /executions /plugin /plugins /build执行mvn clean compile后target/generated-sources/protobuf/java里会有消息实体类target/generated-sources/protobuf/grpc-java里会有MessageServiceGrpc。注意compile-custom这个 goal 负责生成 gRPC stub如果漏掉就只能看到 message 类没有 service 类。还有一个常见坑第一次编译报Could not find any matches for ... as a protocol buffer plugin多半是 protoc 平台检测失败这时确认 os-maven-plugin 是否加到了 build extensions 里。Windows 下如果杀毒软件拦截了 protoc 的 exe 执行也可能报类似错误把生成的临时目录加白名单就好。3. 服务端核心实现3.1 创建 GrpcServer 并注册服务Vert.x 里启动一个 gRPC 服务端非常直接。用GrpcServer.server(vertx)创建然后把实现了MessageServiceImplBase的实例 addService最后调用 listen 绑定端口。这是一个异步方法返回 Future建议挂在 onSuccess 回调里打日志。Vertx vertx Vertx.vertx(); GrpcServer server GrpcServer.server(vertx); MessageServiceGrpc.MessageServiceImplBase serviceImpl new MessageServiceGrpc.MessageServiceImplBase() { Override public void sendMessage(SendRequest request, StreamObserverSendReply responseObserver) { boolean ok messageSender.deliver(request); SendReply reply SendReply.newBuilder() .setSuccess(ok) .setMessageId(request.getMessageId()) .setMessage(ok ? sent : failed) .build(); responseObserver.onNext(reply); responseObserver.onCompleted(); } Override public StreamObserverSendRequest sendMessageStream(StreamObserverSendReply responseObserver) { return new StreamObserverSendRequest() { private int received 0; private SendReply.Builder batchReply SendReply.newBuilder().setSuccess(true); Override public void onNext(SendRequest request) { received; messageSender.deliver(request); } Override public void onError(Throwable t) { log.error(stream error, t); } Override public void onCompleted() { SendReply reply batchReply .setMessage(received received messages) .build(); responseObserver.onNext(reply); responseObserver.onCompleted(); } }; } }; server.addService(serviceImpl) .listen(8080, 0.0.0.0) .onSuccess(v - log.info(gRPC server listening on 8080));一元消息的处理逻辑很好理解收到 request 后做投递然后 onNext 返回结果最后必须 onCompleted。这里要强调一个 gRPC 的约定一次 RPC 调用里onNext最多调用一次在一元模式下并且最终必须要调用onCompleted否则客户端会一直挂着等结果。双向流的实现则返回一个StreamObserverSendRequestgRPC 框架会通过这个 observer 把客户端发来的请求逐个推给我们。我在 onNext 里计数和逐条投递客户端结束发送后会触发 onCompleted服务端再把汇总结果返回给对方。这种模式特别适合批量消息上报。3.2 投递逻辑与异步化处理服务端收到消息后的投递动作可能是写数据库、调第三方接口、写入消息队列这些操作都不应该直接放在 gRPC 回调线程上同步执行太久。我这里的messageSender.deliver内部用 Vert.x 的 Future 链串起来先做鉴权再发到 MQ最后更新发送状态全程异步。public FutureBoolean deliver(SendRequest request) { return authService.check(request.getFrom()) .compose(ok - mqClient.publish(msg.send, request.toByteArray())) .map(sent - true) .otherwise(err - { log.error(deliver failed, err); return false; }); }回调线程是 gRPC 的工作线程如果你塞入了阻塞操作就会把线程池拖住。Vert.x 的妙处在于把这种异步串成一条链代码看起来像同步的实际跑在事件循环上。有个原则我始终遵守回调里只做轻量计算和 responseObserver 回复重活全部丢给 Future 链。3.3 服务端常用参数与大小限制gRPC 默认的消息大小限制是 4MB在 Vert.x 里有方法可以直接调整。GrpcServerOptions提供了setMaxMessageSize单位是字节。我做消息发送时服务端单条最大消息设到 16MB因为有些业务会内嵌 base64 图片4MB 根本不够。GrpcServerOptions options new GrpcServerOptions() .setMaxMessageSize(16 * 1024 * 1024) .setMaxInboundMessageSize(16 * 1024 * 1024) .setIdleTimeout(30) .setIdleTimeoutUnit(TimeUnit.SECONDS); GrpcServer server GrpcServer.server(vertx, options);setMaxInboundMessageSize是接收请求的最大值setMaxMessageSize更像是对外的总限额。两者最好都调否则会出现一边放行一边拦截的诡异情况。客户端默认的发送限额也要对应调整如果客户端不放开写入限额服务端调得再大也收不到大消息。4. 客户端与消息发送实操4.1 创建 GrpcClientChannel 与 Stub 复用客户端的核心是 Channel。Vert.x 的GrpcClient通过GrpcClientChannel包装成 gRPC 能识别的 Channel然后可以用这个 Channel 创建各种 stub。GrpcClient grpcClient GrpcClient.client(vertx); GrpcClientChannel channel new GrpcClientChannel( grpcClient, new VertxChannelFactory(vertx, new NettyClientOptions())); MessageServiceGrpc.MessageServiceStub stub MessageServiceGrpc.newStub(channel);这个 channel 是重量级对象里面维护了连接、线程分配、缓冲区必须复用。我是在 Verticle 的 start 方法里初始化存储成一个字段服务停止时再关闭。关闭逻辑不能漏如果不关闭 grpcClient即使 vertx 实例关闭了底层 gRPC 的线程池也可能会悬挂进程导致优雅停机后进程退不出去这是我在 Vert.x 4.x 版本实测踩过的坑。当初排查这个问题花了不少时间最后用 shutdown hook 主动关闭了 client进程才能正常退出。Override public void stop() { if (grpcClient ! null) { grpcClient.close().onComplete(ar - log.info(gRPC client closed)); } }如果你在多个 Verticle 里都创建了 GrpcClient最好用一个共享的工具类统一管理避免每个 Verticle 各建一套连接白白消耗文件描述符。4.2 一元消息发送与结果确认一元消息最简单构建 request调用 stub 发送在回调里处理结果。SendRequest request SendRequest.newBuilder() .setMessageId(UUID.randomUUID().toString()) .setFrom(system) .setTo(recipient) .setContent(payload) .setTimestamp(System.currentTimeMillis()) .build(); stub.sendMessage(request, new StreamObserverSendReply() { Override public void onNext(SendReply reply) { if (reply.getSuccess()) { log.info(send ok, messageId{}, traceId{}, reply.getMessageId(), reply.getTraceId()); } else { log.warn(send failed, reason{}, reply.getMessage()); } } Override public void onError(Throwable t) { log.error(rpc error, t); } Override public void onCompleted() { log.debug(rpc completed); } });这里有个非常关键的点onNext返回的只是服务端受理后的应答它不代表接收方真的看到了消息。“发送成功”在业务上必须分成两跳第一跳是服务端收到并写入存储第二跳是接收方消费。如果你的业务需要确认接收方真正处理了就要在消息内容里带业务 ID由接收方再反向调一个 ack 接口。我在对接企业微信应用消息时就遇到过这个问题光靠发送接口的返回码不够还要再查发送结果才能判断用户到底有没有收到。4.3 双向流式消息发送的两种实现双向流是消息发送场景里最有价值的部分。客户端可以持续往服务端灌消息服务端也可以持续回结果。Vert.x 的异步 stub 用起来是这样的StreamObserverSendReply replyObserver new StreamObserverSendReply() { Override public void onNext(SendReply reply) { log.info(server reply: {}, reply.getMessage()); } Override public void onError(Throwable t) { log.error(stream error, t); } Override public void onCompleted() { log.info(stream finished); } }; StreamObserverSendRequest requestObserver stub.sendMessageStream(replyObserver); for (int i 0; i 1000; i) { SendRequest request SendRequest.newBuilder() .setMessageId(String.valueOf(i)) .setFrom(producer) .setTo(consumer) .setContent(message- i) .setTimestamp(System.currentTimeMillis()) .build(); requestObserver.onNext(request); } requestObserver.onCompleted();注意一点不要在 for 循环里连续调用 onNext 而不顾对端处理速度。虽然是异步协议但你也可能把内存里待发送的队列撑爆。安全做法是判断一下 write queue 的状态或者用 Vert.x 的vertx.setPeriodic做带节奏的生产每 50ms 发一批把峰值平滑掉。如果是日志类场景更建议用客户端流式一次性发送完再 onCompleted让服务端批量处理。4.4 超时、重试与退避策略gRPC 客户端默认超时是没有上限的这对生产环境很危险。我用Deadline给每个调用设置了 5 秒超时超过就触发DEADLINE_EXCEEDED。stub.withDeadlineAfter(5, TimeUnit.SECONDS).sendMessage(request, observer);更细一点可以用StreamObserver.onError捕获异常码根据状态码决定重试策略。我一般这样处理gRPC 状态码含义是否重试UNAVAILABLE服务不可用是用指数退避DEADLINE_EXCEEDED超时是但退避时间要放大INTERNAL服务端内部错误否大概率是参数或服务 bugINVALID_ARGUMENT请求参数非法否重试也是白搭RESOURCE_EXHAUSTED限流或资源耗尽是但必须等待更长时间重试不能全量重发否则服务端刚缓过来又被你打挂。退避我推荐从 200ms 开始每次乘 2最多到 5 秒。这是基于实际踩过限流问题的教训——消息发送过于频繁被服务端拒绝时代码里如果没有退避客户端会在秒级内疯狂重试等于给自己制造雪崩。5. 常见问题与排查技巧实录5.1 依赖冲突两类 Netty 的战争Vert.x 内置了 Netty 4.1grpc-netty 也带了 Netty 4.1版本不一致时 classpath 里会出现大量netty-tcnative、netty-handler之类的 jar。报错通常很惊悚比如NoSuchMethodError或ClassNotFoundException。解决方案我推荐优先使用grpc-netty-shaded它的 Netty 类被打包到了io.grpc.netty.shaded包下和 Vert.x 的 Netty 互不干扰。如果必须用grpc-netty就要把复杂依赖版本统一到一致常见组合是 Netty 4.1.100 和 grpc-java 1.62但这很折腾不值得。5.2 连接泄漏与优雅停机另一个高频问题是GrpcClientChannel和GrpcClient没有关闭。跑一段时间后连接数只增不减到最后报OutOfMemoryError: unable to create new native thread。我排查时用netstat -anp | grep 8080看连接数发现 TIME_WAIT 和 ESTABLISHED 都在涨。为什么涨因为很多人只创建了 stub没持有 channel 的引用导致 channel 被垃圾回收后底层连接没有显式关闭gRPC 后台线程池又一直挂着。解决很简单channel 用单例持有应用退出时调用 grpcClient.close()而这个 close 在 Vert.x 里是个异步操作要等待它的 Future 完成再销毁 Vertx。5.3 大消息发送失败与性能调优默认 4MB 的限制会让业务发大 JSON 时报RESOURCE_EXHAUSTEDgRPC 日志提示message too large。除了服务端要把maxMessageSize调大客户端也必须用同样的方式放开否则服务端收不到完整数据。传输大消息时 protobuf 的性能优势会体现得比小消息更明显因为它避免了 JSON 反复解析的 CPU 开销。我在压测里试过一个 8MB 的 base64 字段REST 的发送耗时大概是 gRPC 的 1.6 倍吞吐也差不少。但大消息会占住内存建议对超大内容做分片或压缩在 proto 里增加 compression 字段用 gzip 压缩 content实测对文本类消息能有 5 到 10 倍的体积缩减。5.4 如何有效确认消息发送成功最后回到很多人会问的确认问题。gRPC 的 reply 只能说明服务端受理成功不保证业务送达。我的方案是三层确认第一层是 RPC 层的 reply用于确认网络传输和服务端处理无异常。第二层是服务端返回的 traceId客户端拿它去查询服务端日志确认消息写入了目标队列或存储。第三层是接收方回调适用于强一致场景接收方处理完业务后调用一个 ack 接口发送方收到 ack 才真正认为消息发送成功。在代码里我会把第一层和第二层合并成一句话reply.getSuccess() StringUtils.hasText(reply.getTraceId())作为最基本的发送成功判断。如果这个都不满足直接走异常流程或者补发。这个判断简单、直观排障时也容易对齐日志关键字。到这里基于 Vert.x 的 gRPC 消息发送从设计、配置到实现、排查的主线就完整了。我个人最想留给你的一句话是gRPC 看起来很重其实它把通信层的复杂度处理得很干净真正决定一个消息系统稳不稳的是你对异步模型、超时重试、资源回收这三个基础动作的态度。先把这几个动作做对比研究任何高级特性都有用。