资讯动态

Vert.x gRPC实践:构建高性能消息发送服务指南

发布时间:2026/10/5 3:07:21 来源:尧图企业网站定制
1. 为什么用 Vert.x gRPC 组合来做消息发送1.1 从一次消息推送需求说起先交代一下背景。前阵子接了一个内部系统改造要求把原来基于 HTTP 轮询拿数据的逻辑换成服务端主动推送消息而且推送的格式、字段、版本都要严格统一调用方有 Java 客户端、有 Python 脚本、还有跑在手机上的 App 端。第一反应是直接上 gRPC因为 Protocol Buffers 定义接口和消息结构的能力太强了跨语言兼容几乎没有成本。但问题是gRPC 原生的 Java 实现走的是 Netty而我们的核心服务跑在 Vert.x 上需要的是那套 Event Loop 线程模型、异步调链、以及一套代码同时处理 HTTP/TCP/消息队列的能力。既然都是 Netty 底层为什么不直接让 gRPC 融入到 Vert.x 的异步体系里呢这就是 Vert.x gRPC 的典型使用场景既想要 gRPC 带来的强契约、二进制序列化、双向流式通信能力又不想放弃 Vert.x 的响应式编程模型和资源管理方式。对做消息发送这类高频、短连接、强确认诉求的场景这个组合其实比很多人想象中要舒服得多。1.2 gRPC 比 REST 和消息队列好在哪聊 gRPC 做消息发送难免会有人问为什么不直接 HTTP Post为什么不用 Kafka 或者 RabbitMQHTTP REST 最大的问题不是性能而是没有一个统一、强约束的接口契约。你用 JSON 传消息字段类型、嵌套结构、是否必填基本都是靠“君子协定”客户端和服务端一不一致全靠 Code Review。消息发送这种场景最怕什么怕线上字段对不上怕服务端改了字段名客户端不知道。gRPC 用 .proto 文件同时生成服务端和客户端代码接口契约在编译期就锁死了字段增删有编号规则兼容性有官方背书这对消息类接口来说是决定性的优势。至于消息队列它的定位是“削峰填谷、异步解耦”而 Vert.x gRPC 解决的是“点对点实时消息投递与确认”。两者不是替代关系很多系统里是同时存在的外部请求用 gRPC 进来落到 Vert.x 的 Event Bus 里做异步分发再落到 MQ 做持久化和重试。单论“消息发送”这个动作本身gRPC 的 unary 调用和双向流都能给出实时性确定性很强的行为这是 MQ 天然不擅长的。1.3 Vert.x 异步模型对 gRPC 的无缝适配很多人在传统 Java gRPC 里体验不好根源在于回调线程模型gRPC 的 Listener 回调跑在 Netty 的 EventLoop 线程上你的业务回调一旦做了阻塞操作整个 Netty worker 线程就被拖住了吞吐量直接垮掉。Vert.x 不一样。Vert.x 从设计上就强制开发者走 Event Loop 回调/协程的异步模式它不允许你在 Event Loop 上做阻塞操作所有耗时逻辑要么扔到 Worker 线程池要么写成异步调用。Vert.x 的 gRPC 实现把 gRPC 的回调接入到了 Vert.x 的 Context 体系里你在回调里拿到结果后可以非常自然地继续用 Vert.x 的 Future、RxJava 或者 Kotlin 协程编排后续逻辑不会破坏调用链。直接说结论传统 gRPC 适合“接口调用是孤立的、不走统一线程模型”的场景而 Vert.x gRPC 适合“消息发送是整条响应式调用链上的一环”的场景。我们的项目里消息进来后要经过鉴权、过滤、持久化、限流、风控等多个环节用 Vert.x 的 Future 把 gRPC 调用串进去简直顺手得不行。这也是整个技术选型最核心的决策点。2. 环境准备与依赖落地2.1 版本选型Vert.x 4.x gRPC 生态做 Vert.x 的 gRPC 开发第一步就是版本匹配这一步踩坑概率最高。Vert.x gRPC 模块在 4.x 时代已经比较成熟官方把vertx-grpc拆成了服务端和客户端两组依赖底层通信协议库用的还是 gRPC Java 那套只不过把传输层和 Vert.x 做了绑定。我自己测试下来比较顺的版本组合是Vert.x 4.5.x gRPC Java 1.6x Protobuf Java 3.2x。注意 Vert.x 5.x 也有对应的 gRPC 模块但如果你在生产环境跑的是 4.x 的老项目不建议为了 gRPC 单独升大版本集成成本和回归风险都偏高。用 Maven 的话推荐直接用 BOM 管理 Vert.x 版本避免多个 Vert.x 模块版本不一致导致 NoClassDefFoundError。2.2 Maven 依赖配置样例直接上配置。我用的是 Maven 多模块工程协议定义单独放在一个api模块里负责生成 Java 代码服务端和客户端各自依赖这个模块。properties vertx.version4.5.10/vertx.version grpc.version1.64.0/grpc.version protobuf.version3.25.3/protobuf.version /properties dependencyManagement dependencies dependency groupIdio.vertx/groupId artifactIdvertx-stack-depchain/artifactId version${vertx.version}/version typepom/type scopeimport/scope /dependency /dependencies /dependencyManagement核心依赖!-- Vert.x gRPC 服务端与客户端 -- dependency groupIdio.vertx/groupId artifactIdvertx-grpc/artifactId /dependency !-- gRPC 生态基础库 -- 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 groupIdjavax.annotation/groupId artifactIdjavax.annotation-api/artifactId version1.3.2/version /dependency !-- protobuf 编译插件 -- plugin groupIdcom.github.os72/groupId artifactIdprotoc-jar-maven-plugin/artifactId version3.11.4/version executions execution phasegenerate-sources/phase goals goalrun/goal /goals configuration includeMavenTypesdirect/includeMavenTypes inputDirectories includesrc/main/proto/include /inputDirectories outputTargets outputTarget typejava/type outputDirectorysrc/main/java/outputDirectory /outputTarget /outputTargets /configuration /execution /executions /plugin有个细节值得提醒grpc-netty-shaded和vertx-grpc一起用时不要再额外引入原生的grpc-netty否则会出现类冲突。Vert.x 的 gRPC 服务端其实是基于自己的 Netty 实例创建的并不需要额外起 gRPC 的 HTTP/2 服务。依赖引错的情况下最常见的报错是ClassNotFoundException: io.grpc.internal.ServerImpl或者NoSuchMethodError这种问题排查起来特别费时间细心检查依赖树是第一步。3. 核心实现proto 定义与服务端开发3.1 先把消息协议定义到 proto 文件里消息发送的第一件事是定义清楚服务接口和消息体。这里拿一个极简而贴近实际的消息推送服务举例这个服务要支持单条消息发送和批量消息发送两种操作并且要明确返回每条消息的发送结果。syntax proto3; package message.gateway; option java_multiple_files true; option java_package com.example.message.gateway; option java_outer_classname MessageGatewayProto; service MessageGateway { // 单条消息发送 rpc SendMessage(SendMessageRequest) returns (SendMessageResponse); // 批量消息发送 rpc BatchSendMessages(BatchSendMessagesRequest) returns (stream SendMessageResponse); } message SendMessageRequest { string message_id 1; string channel 2; string receiver 3; string content 4; mapstring, string headers 5; } message BatchSendMessagesRequest { repeated SendMessageRequest messages 1; } message SendMessageResponse { string message_id 1; bool success 2; string code 3; string message 4; }对这个 proto 有几点解读option java_multiple_files true让每个 message 生成独立的 Java 类避免一个巨大的外层类代码组织更清晰。字段编号一旦确定就不要改。消息发送牵涉到历史数据、调用方新旧版本并存如果改了字段编号老客户端反序列化时会出现字段错位这个坑在 protobuf 里极其隐蔽。单个 rpc 返回的消息体里带上message_id是为了做幂等。消息发送场景最怕客户端超时重试造成重复发送服务端拿到 message_id 可以在内存或 Redis 里做去重。3.2 服务端实现把业务逻辑绑定到 Vert.x 生命周期Vert.x gRPC 服务端的核心入口是GrpcServer它可以在已有的Vertx实例上创建。所有 gRPC 调用会像 Vert.x 的普通请求一样分发到 Event Loop 上执行。我的实现方式是这样的public class MessageGatewayServer { public static void main(String[] args) { Vertx vertx Vertx.vertx(); // 创建 gRPC 服务端监听 8080 端口 GrpcServer server GrpcServer.create(vertx, new GrpcServerOptions() .setHost(0.0.0.0) .setPort(8080)); // 实例化业务处理类 MessageGatewayService service new MessageGatewayService(vertx); // 将 gRPC 服务绑定到 Vert.x 服务端 server.callHandler(MessageGatewayGrpc.getSendMessageMethod(), new VertxMessageGatewayImplBase(service) { Override public FutureSendMessageResponse sendMessage(SendMessageRequest request) { return service.handleSend(request); } }); server.start().onComplete(ar - { if (ar.succeeded()) { System.out.println(gRPC server started on port 8080); } else { System.err.println(gRPC server failed to start); } }); } }需要注意VertxMessageGatewayImplBase是 Vert.x gRPC 针对你的 proto 自动生成的抽象基类和原生 gRPC 生成的MessageGatewayGrpc.MessageGatewayImplBase不同它返回的是 Vert.x 风格的FutureT不是在回调里手动调用responseObserver。这样一来业务逻辑里可以直接使用io.vertx.core.Future的链式调用和 Vert.x 生态完全打通。服务端的handleSend方法举例如下public class MessageGatewayService { private final Vertx vertx; private final MessageStore store; public MessageGatewayService(Vertx vertx) { this.vertx vertx; this.store new MessageStore(vertx); } public FutureSendMessageResponse handleSend(SendMessageRequest request) { // 1. 先做幂等检查消息 ID 是否已存在 return store.existMessage(request.getMessageId()) .compose(exists - { if (exists) { // 已存在直接返回成功表示重复投递已被去重 return Future.succeededFuture(SendMessageResponse.newBuilder() .setMessageId(request.getMessageId()) .setSuccess(true) .setCode(DUPLICATE_IGNORED) .build()); } // 2. 模拟消息写入存储 投递到下游 return store.saveMessage(request) .compose(v - deliverToReceiver(request)) .map(v - SendMessageResponse.newBuilder() .setMessageId(request.getMessageId()) .setSuccess(true) .setCode(OK) .build()) .otherwise(err - SendMessageResponse.newBuilder() .setMessageId(request.getMessageId()) .setSuccess(false) .setCode(DELIVERY_FAILED) .setMessage(err.getMessage()) .build()); }); } }这里我特别强调一下 Event Loop 纪律handleSend里不要出现任何阻塞调用。如果store.saveMessage是查 MySQL 或者 Redis一定要用客户端自己的异步 API比如vertx-mysql-client、vertx-redis-client而不是在 Event Loop 上直接跑 JDBC。实际项目里见过有人为省事在 gRPC 回调里写 JDBC结果数据库一慢整个 Verticle 的 Event Loop 全被拖死所有请求排队这是 Vert.x 开发中比较典型的反面教材。3.3 批量发送处理客户端流还是要关注响应顺序批量发送“消息”这个场景在 proto 里我特意用了服务端流式返回客户端一次提交多条消息服务端逐条处理并返回结果。这种模式适合发通知、发营销消息这类量大但每条之间相互独立的业务逻辑。public void batchSendMessages(BatchSendMessagesRequest request, PromiseSendMessageResponse response) { ListSendMessageRequest messages request.getMessagesList(); for (SendMessageRequest msg : messages) { // 每条消息独立处理用 map 拼结果保证顺序与输入一致 handleSend(msg).onComplete(ar - { if (ar.succeeded()) { response.complete(ar.result()); } else { response.fail(ar.cause()); } }); } }当然更稳妥的写法是每条响应通过Promise逐个处理而 gRPC 的服务端流式接口要求客户端订阅FlowableSendMessageResponse来接收多条响应。这里有一个很现实的业务约定要提前想清楚批量的消息是逐条实时返回还是等全部处理完成一次性返回如果是逐条返回用服务端流式很自然如果下游需要所有消息都到达后才处理建议 unary 返回一个BatchSendMessagesResponse包装列表否则客户端处理 partial result 的逻辑会很别扭。4. 客户端实现与消息发送的异步化4.1 用异步 Stub 而不是 Blocking Stub客户端同样有原生 gRPC 的blockingStub和asyncStub两种用法但在 Vert.x 环境里官方推荐使用基于 Vert.x 包装的异步 Stub它在每次调用时返回FutureSendMessageResponse可以直接和项目的异步调用链对接。手动调用下面的代码完成客户端编写public class MessageGatewayClient { private final MessageGatewayVertxGrpc.MessageGatewayVertxStub stub; public MessageGatewayClient(Vertx vertx, String host, int port) { GrpcClient client GrpcClient.create(vertx, new GrpcClientOptions()); this.stub MessageGatewayVertxGrpc.newStub(client, new GrpcClientRequestOptions() .setHost(host) .setPort(port)); } public FutureSendMessageResponse send(SendMessageRequest request) { return stub.sendMessage(request); } }这里有个细节GrpcClient可以复用。客户端初始化一次后续所有消息发送都共用同一个 gRPC 连接这得益于 HTTP/2 的多路复用特性同一个连接上可以同时跑大量并发请求不会像 HTTP/1.1 那样需要频繁建立和断开 TCP 连接。4.2 超时设置、失败确认与重试注意事项消息发送能不能确定成功是业务上最关心的事。gRPC 天然提供了三层的“确认信号”调用正常返回SendMessageResponse.success true业务成功。调用返回业务失败success false但连接层面正常。调用抛异常比如StatusRuntimeException说明请求可能都没到业务层或者服务端处理时挂了。这种场景最需要小心因为客户端无法判断服务端是否已经收到并开始处理了贸然重试可能造成重复消息。这时候message_id的幂等保护就很重要。超时配置上我建议在客户端为每个调用设置 deadline避免下游服务挂死时客户端无限等待。gRPC 的 deadline 与业务超时不同它走的是 HTTP/2 的 ping 机制比较直接。public FutureSendMessageResponse sendWithTimeout(SendMessageRequest request) { // 在客户端把超时时间传入 stub 的 deadline return stub.sendMessage(request); }实际用的时候如果希望整体控制超时可以统一通过GrpcClientRequestOptions设置 deadline。另外说一点客户端重试不是默认开启的需要显式配置RetryPolicy。消息发送场景我会建议“幂等 有限重试”即调用失败后最多重试两次并且两次重试之间要有指数退避的间隔否则在服务端已经过载的情况下重试只会加重问题。4.3 发送过于频繁时的限流策略架构上想清楚“消息发送过于频繁请稍后重试”这类报错本质上是一种保护机制服务方主动限制客户端发送速率。实现上有两种思路服务端限流和客户端限流。服务端限流一般用令牌桶对每个发起方或者每个 channel 维度设置 QPS 上限。判断超限后直接返回一个可重试的错误码。客户端这边Vert.x 配合 Guava 的 RateLimiter 也能实现简单的单机限流比如限制每秒最多发送 N 条public class ThrottledMessageGatewayClient { private final MessageGatewayVertxGrpc.MessageGatewayVertxStub stub; private final RateLimiter rateLimiter RateLimiter.create(100.0); // 每秒最多 100 条 public FutureSendMessageResponse send(SendMessageRequest request) { // 若令牌不足可以阻塞等待也可以立刻返回限流异常看业务取舍 if (!rateLimiter.tryAcquire(1)) { return Future.failedFuture(too many requests, please retry later); } return stub.sendMessage(request); } }这里选择“立刻失败”还是“等待令牌”取决于调用方的忍耐度。如果是内部系统一般我会建议用一个带缓冲的有界队列接不住就直接失败并反馈重试提示不要在生产环境里让客户端无限阻塞下去不然会把线程池打满。类似“deepseek消息发送过于频繁请稍后重试”这类提示本质就是上游在做限流保护客户端要做的不是硬扛而是合理退避。5. 实际部署中踩过的坑与问题排查5.1 连接管理频繁创建客户端导致端口耗尽第一次压测的时候遇到过一个特别诡异的故障客户端高并发跑了几万条消息后突然大量请求超时服务端日志显示连接被重置。排查半天发现问题出在测试代码里每次发送都new一个GrpcClient。每个GrpcClient会维护独立的 HTTP/2 连接和线程池创建太多后文件描述符和本地端口全部被耗尽连接根本建不出来。正确做法是全局只建一个GrpcClient实例通过GrpcClientRequestOptions来区分不同的目标地址。如果一个客户端要连多个服务端也是同一个GrpcClient负责全部连接不必重复创建。5.2 消息体过大导致服务端报错默认情况下 gRPC 服务端收到的消息体大小限制是 4MB客户端可接收的响应大小限制也是 4MB。消息发送场景里如果 content 字段里塞了一大段日志或者 Base64 图片很容易触发RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 4194304 bytes解决方案是在服务端和客户端都显式调大限制值// 服务端 GrpcServerOptions options new GrpcServerOptions(); options.setMaxInboundMessageSize(32 * 1024 * 1024); // 32MB // 客户端 GrpcClientOptions clientOptions new GrpcClientOptions(); clientOptions.setMaxInboundMessageSize(32 * 1024 * 1024);但我要多说一句调大限制只是解表不治本。消息体过大会导致序列化耗时、GC 压力、带宽占用全线飙升业务上最好还是拆分消息不要让单条消息承载太多数据。5.3 确认发送成功的排查思路有朋友问企业微信发送应用消息是怎么确认发送成功的这个问题其实和 gRPC 的场景类似。应用消息发送的确认可以分为两层第一层是“接口调用成功”第二层是“用户真正收到了”。gRPC 的 unary response 给你的是第一层确认服务端业务代码里确认消息已落库、已投递到下游入口才返回 success。如果你还需要第二层确认那就要在协议里额外设计回执字段或者回调通知。排查“消息没发送成功”时我一般会按这几个步骤逐步推进服务端有没有收到请求在服务端入口打日志记录 message_id 和 channel。gRPC 调用返回什么状态如果Status.UNIMPLEMENTED大概率是方法路径对不上如果是Status.DEADLINE_EXCEEDED优先查服务端有没有阻塞操作。消息落库成功了吗如果落库失败而客户端拿到了成功返回一定是业务代码里 try-catch 吞掉了异常这是比较低级但很常见的问题。客户端确认成功的语义是不是透传了上游状态不要把“请求发出去了”当成“发送成功”除非你用的是 Fire-and-Forget。5.4 常见问题速查表现象可能原因排查方法调用一直超时服务端线程被阻塞Event Loop 打满打开 Vert.x 的线程体检vertx-meter或检查日志里的 blocked thread客户端连接被重置创建了过多 GrpcClient 连接改为单例复用检查文件描述符消息超过最大限制默认 4MB 限制服务端和客户端同时调大setMaxInboundMessageSize或者拆包随机的重复消息客户端超时重试但服务端实际已处理用 message_id 做幂等去重批量发送响应乱序服务端多条 Promise 异步完成顺序不一致保证处理循环顺序提交 Promise或收集后统一返回压测时性能下降客户端没有启用连接池复用复用 GrpcClient确认 HTTP/2 多路复用生效proto 新增字段后老客户段解析出错字段编号被复用或删除proto 字段编号一经定义终身保留删除用 reserved 声明5.5 关于 gRPC 和 MQ 的边界再补一刀虽然前面说了 gRPC 和 MQ 不冲突但这里还是要给一些经验不足的读者提个醒如果你的消息发送场景长期来看需要可靠投递、消息堆积、消费失败重试、广播等能力那 gRPC 顶多只能作为入口层最终还是要落到 MQ 上。gRPC 适合做同步、低时延、需要立刻知道消息“发出去了没有”的场景MQ 适合做需要容忍秒级延迟、支持大量堆积与回溯的场景。两个不是谁替代谁而是上下游配合的关系。把这点想清楚技术选型上就不会再纠结。6. 用 Vert.x gRPC 做消息发送的更多扩展思路前面讲的是基础玩法实际上消息发送在真实系统里往往还有很多变种这里我挑两个比较常见的扩展方向给大家一个参考。6.1 双向流式消息下发单向 unary 适合“请求 - 响应”的简单消息发送但如果你的业务是“客户端建立连接后服务端持续向客户端推送消息”这就非常适合使用 gRPC 的双向流。Vert.x 的 gRPC 对双向流的支持非常友好客户端和服务端可以持续写入多条消息形成类似 websocket 的效果但又有 protobuf 的强类型约束。// 服务端定义 rpc ChannelStream(stream ChannelRequest) returns (stream PushMessage);在这个模型下客户端可以长期维持一个连接服务端向客户端发送实时消息比如行情推送、站内信、设备指令下发都不需要客户端反复发起请求。对比 MQTT 或者原生 WebSocketgRPC 双工流在微服务互联场景下的接入成本和运维规范性更好尤其是服务端之间需要程序化通信时比 WebSocket 更顺手。6.2 消息转发到事件总线Vert.x 内置了EventBus这是它非常强大的一个能力。很多推送场景不需要 gRPC 服务端直接处理所有业务而是把收到的消息编码后发到 EventBus 上由其他 Verticle 消费处理。public FutureSendMessageResponse handleSend(SendMessageRequest request) { // 把消息发布到 event bus 上让其他消费者异步处理 return vertx.eventBus().request(msg.gateway.send, JsonObject.mapFrom(request)) .map(msg - SendMessageResponse.newBuilder() .setMessageId(request.getMessageId()) .setSuccess(true) .setCode(ACCEPTED) .build()); }这样做的好处是gRPC 层只管接入和响应具体投递逻辑由下游 Verticle 负责模块间耦合更低水平扩展也更灵活。消息量上来以后还可以给 EventBus 添加集群模式多个 Vert.x 节点共享同一套消息地址发送能力可以横向扩展。7. 实操后的心得与建议说实话Vert.x gRPC 这套组合在国内技术社区里的讨论热度一直不如 Spring Boot gRPC 或者 Spring Cloud 那套高但它其实非常适合已经上了 Vert.x 船、又被跨语言通信和强契约问题困住的项目。如果你正在评估方案我给几个实际的建议第一proto 文件是核心资产避免随便改。把它当成接口文档来维护字段增删要评审编号绝不重复使用。建议用单独仓库管理 proto配合 Buf 这类 lint 工具做风格检查和兼容性校验。第二客户端链路尽早把超时和重试策略定下来。消息发送这类操作用户感知最强的是“到底发出去没有”如果超时策略混乱重试逻辑复杂线上排查会非常痛苦。先约定好哪些错误码是可重试的哪些是不能重试的再写代码。第三利用好 Vert.x 的调试手段。开发时开启vertx.logger-delegate-factory-class-name和 gRPC 的日志等级能看到底层的 HTTP/2 帧交互排查连接问题是利器。生产环境则要结合 Micrometer 之类的监控系统把 gRPC 的请求耗时、错误率、消息体大小都曝露出来否则出了问题只能靠猜。最后再分享一个小技巧在压测消息发送接口时不要只看 QPS还要关注服务端 Event Loop 的延迟和线程占用。gRPC 这层如果阻塞了 Event Loop哪怕 QPS 不高整体响应也会像蜗牛一样慢。提前在测试环境压出问题优于线上告警之后再救火。消息发送的本质是交付和确认把这条链路的每一个环节都量化好这套方案才能跑得又稳又长久。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑