gRPC Java如何获取流读取低级事件或实现流读取器拦截
gRPC 流读取低层事件拦截方案
默认的StreamObserver确实仅会在完整消息反序列化完成后触发onNext回调,你需要的「流读取进行中」类低层级事件、流读取器拦截能力,均可通过gRPC官方提供的扩展机制实现:
可用实现方案
方案1:通过拦截器自定义流监听器
这是适配性最强的方案,客户端、服务端均可使用:
- 客户端侧实现
ClientInterceptor接口,重写interceptCall方法,替换默认的ClientCall.Listener,即可在流读取的各个阶段拿到回调:onReady():流就绪可接收数据时触发onMessage(InputStream message):每收到一段数据分片就会触发,此时数据还未完成反序列化,你可以在这里重置应用层心跳的超时计时器,避免大消息传输过程中误判超时- 处理完自定义逻辑后再把事件转发给原有
StreamObserver即可不影响原有业务逻辑
- 服务端侧同理实现
ServerInterceptor接口,重写interceptCall方法替换ServerCall.Listener,即可拿到服务端侧的流读取各阶段事件
方案2:自定义消息解析器
如果你不需要太细的网络层事件,仅需要感知消息解析进度,可以自定义实现Marshaller接口:
- 重写
parse(InputStream input)方法,在读取输入流反序列化的过程中,每读取一段数据就触发自定义的进度回调即可,不需要修改拦截器逻辑,仅需要在生成gRPC stub的时候指定自定义的Marshaller即可生效
针对你的场景的优化建议
你当前遇到的10秒应用层ping超时问题,可以直接通过上述方案解决:
- 只要在收到任何流数据分片的时候就重置超时计时器,不需要等完整消息反序列化完成触发
onNext再重置,即可避免慢速网络下大消息传输导致的误超时 - 如果不需要精确的读取进度,仅需要保活的话,也可以直接开启gRPC内置的
keepalive机制,配置keepalive_time小于你的应用层超时时间,gRPC内核会自动发送传输层ping帧,不需要应用层额外处理大消息传输的保活逻辑
相关接口中文翻译
public interface StreamObserver<V> { /** * 接收流中的值 * * <p>可被多次调用,但在{@link #onError(Throwable)}或{@link #onCompleted()}调用后不会再被调用 * * <p>一元调用最多触发一次onNext。服务端流调用场景下客户端最多调用一次onNext,但可以接收多次onNext回调;客户端流调用场景下服务端最多调用一次onNext,但可以接收多次onNext回调 * * <p>如果实现类抛出异常,调用方需要先通过{@link #onError(Throwable)}传入捕获的异常终止流,再传播该异常 * * @param value 传入流的值 */ void onNext(V value); /** * 接收流的终止错误 * * <p>仅可被调用一次,且必须是最后一个被调用的方法。尤其如果{@code onError}的实现类抛出异常,不允许再调用任何其他方法 * * <p>{@code t}通常是{@link io.grpc.StatusException}或{@link io.grpc.StatusRuntimeException},但也可能是其他{@code Throwable}类型。调用方通常应该通过{@link io.grpc.Status#asException()}或{@link io.grpc.Status#asRuntimeException()}从{@code Status}转换得到;实现类通常应该通过{@link io.grpc.Status#fromThrowable(Throwable)}转换得到{@code Status} * * @param t 流中发生的错误 */ void onError(Throwable t); /** * 接收流成功完成的通知 * * <p>仅可被调用一次,且必须是最后一个被调用的方法。尤其如果{@code onCompleted}的实现类抛出异常,不允许再调用任何其他方法 */ void onCompleted(); }
内容的提问来源于stack exchange,提问作者Harshit Bangar
相关产品推荐
相关产品推荐

