Node.js GRPC客户端流消息未实时发送问题求助
问题:Node.js gRPC客户端流消息未即时发送
我正尝试以Node.js作为客户端(服务端基于.NET)实现gRPC客户端流。以下是proto文件:
syntax = "proto3"; package groom; import "google/protobuf/timestamp.proto"; message NewsFlash { google.protobuf.Timestamp news_time=1; string news_item=2; } message NewsStreamStatus { bool success=1; } service Groom { rpc SendNewsFlash(stream NewsFlash) returns (NewsStreamStatus); }
我的客户端代码如下:
const grpc = require("@grpc/grpc-js"); var protoLoader = require("@grpc/proto-loader"); const PROTO_PATH = "./Protos/groom.proto"; const options = { keepCase: true, longs: String, enums: String, defaults: true, oneofs: true, }; const newsItems=["Item1","Item2","Item3","Item4","Item5"] var grpcObj = protoLoader.loadSync(PROTO_PATH, options); const GroomService = grpc.loadPackageDefinition(grpcObj).groom.Groom; const clientStub = new GroomService( "localhost:5054", grpc.credentials.createInsecure() ); var call = clientStub.sendNewsFlash(function(error, newsStatus) { if (error) { console.error(error); } console.log('Stream success: ', newsStatus.success); }); for (var i=0;i<5;i++) { var itemIndex=Math.floor(Math.random() * 5); call.write({news_item: newsItems[itemIndex]}); } call.end();
代码可正常运行,但消息并未真正流式发送:所有消息仅在调用call.end()时才发送至服务端,而非调用call.write(...)时即时发送。使用BloomRPC模拟调用时服务端能即时接收消息,请问这一问题的原因是什么?
原因及解决方案
核心原因
Node.js的@grpc/grpc-js客户端默认会对发送的消息进行批量缓冲,以此提升传输效率。你的代码是在同步循环中连续调用call.write(),客户端会把所有消息暂存到缓冲区,直到调用call.end()才一次性发送。而BloomRPC是逐条触发消息发送,不会触发批量缓冲逻辑,所以服务端能即时接收。
另外,你发送的NewsFlash消息缺少news_time字段,虽然proto3允许字段缺省,但部分gRPC实现可能会对不完整消息做额外缓冲处理,不过这不是主要原因。
解决方法
要实现即时发送,需要打破同步循环的连续调用,给客户端留出缓冲刷新的时间,同时可以显式控制发送时机:
方案1:异步延迟发送
在每次call.write()后加入短暂延迟,让客户端事件循环有时间处理发送逻辑,同时补充完整的消息字段:
const grpc = require("@grpc/grpc-js"); var protoLoader = require("@grpc/proto-loader"); const PROTO_PATH = "./Protos/groom.proto"; const options = { keepCase: true, longs: String, enums: String, defaults: true, oneofs: true, }; const newsItems=["Item1","Item2","Item3","Item4","Item5"] var grpcObj = protoLoader.loadSync(PROTO_PATH, options); const GroomService = grpc.loadPackageDefinition(grpcObj).groom.Groom; const clientStub = new GroomService( "localhost:5054", grpc.credentials.createInsecure() ); var call = clientStub.sendNewsFlash(function(error, newsStatus) { if (error) { console.error(error); } console.log('Stream success: ', newsStatus.success); }); // 异步逐条发送消息 async function sendMessages() { for (var i=0;i<5;i++) { var itemIndex=Math.floor(Math.random() * 5); const now = new Date(); // 补充完整的news_time字段 call.write({ news_item: newsItems[itemIndex], news_time: { seconds: Math.floor(now.getTime() / 1000), nanos: (now.getTime() % 1000) * 1000000 } }); // 加入短暂延迟,让客户端发送当前消息 await new Promise(resolve => setTimeout(resolve, 100)); } call.end(); } sendMessages();
方案2:监听drain事件精准控制
通过监听call的drain事件,当缓冲区为空时再发送下一条消息,这种方式比固定延迟更高效:
const grpc = require("@grpc/grpc-js"); var protoLoader = require("@grpc/proto-loader"); const PROTO_PATH = "./Protos/groom.proto"; const options = { keepCase: true, longs: String, enums: String, defaults: true, oneofs: true, }; const newsItems=["Item1","Item2","Item3","Item4","Item5"] var grpcObj = protoLoader.loadSync(PROTO_PATH, options); const GroomService = grpc.loadPackageDefinition(grpcObj).groom.Groom; const clientStub = new GroomService( "localhost:5054", grpc.credentials.createInsecure() ); var call = clientStub.sendNewsFlash(function(error, newsStatus) { if (error) { console.error(error); } console.log('Stream success: ', newsStatus.success); }); // 递归发送消息,根据缓冲区状态控制 function sendNextMessage(index) { if (index >= 5) { call.end(); return; } var itemIndex=Math.floor(Math.random() * 5); const now = new Date(); // write方法返回false表示缓冲区已满,需要等待drain事件 const canContinue = call.write({ news_item: newsItems[itemIndex], news_time: { seconds: Math.floor(now.getTime() / 1000), nanos: (now.getTime() % 1000) * 1000000 } }); if (!canContinue) { call.once('drain', () => sendNextMessage(index + 1)); } else { sendNextMessage(index + 1); } } sendNextMessage(0);
内容的提问来源于stack exchange,提问作者ml123
相关产品推荐
相关产品推荐

