Java gRPC双向流异常:客户端取消致网关-游戏服务连接中断
架构背景
正在开发一个依赖Java gRPC双向流的游戏运行时服务器,数据流向:
Client <--grpc bi-stream--> Gateway Server <---- grpc bi-stream -----> Game Server
问题场景
网关服务器维护双向流的Scala代码(JVM兼容,Java开发者可类比):
// Gateway Server gRPC bi-stream API override def stateSyncStreamDemo( responseObserver: StreamObserver[StateSyncFrameDemo] ): StreamObserver[StateSyncStreamRequestDemo] = { // get context from gRPC interceptor val `Context-UserId-Key` = Context.key[String]("user-id") val `Context-RoomId-Key` = Context.key[String]("room-id") new StreamObserver[StateSyncStreamRequestDemo] { override def onNext( request: StateSyncStreamRequestDemo ): Unit = { // send to Game Server by gRPC bi-stream // 1. create bi-stream to Game Server if roomId not create before val streamConnection = getOrCreateRoomConnection(`Context-RoomId-Key`) val gameServerRequestStream = streamConnection.withCallCredentials( new StreamShardingClient.ClientMetadataCall( roomId // setting roomId to Metadata ) ). sendWithBiStream(someResponseStreamDefined) // 2. send message forward to Game Server with userId and client cmd gameServerRequestStream.onNext(Message(request.cmd, userId)) } ... }
该转发模式正常运行,但用grpcurl断开客户端与网关的gRPC连接时,网关与游戏服务器的流也会随之断开,错误日志:
CANCELLED: client cancelled io.grpc.StatusRuntimeException: CANCELLED: client cancelled at io.grpc.Status.asRuntimeException(Status.java:530) ~[grpc-api-1.52.1.jar:1.52.1] at io.grpc.stub.ServerCalls$StreamingServerCallHandler$StreamingServerCallListener.onCancel(ServerCalls.java:291) [grpc-stub-1.52.1.jar:1.52.1] at io.grpc.PartialForwardingServerCallListener.onCancel(PartialForwardingServerCallListener.java:40) [grpc-api-1.52.1.jar:1.52.1] at io.grpc.ForwardingServerCallListener.onCancel(ForwardingServerCallListener.java:23) [grpc-api-1.52.1.jar:1.52.1] at io.grpc.ForwardingServerCallListener$SimpleForwardingServerCallListener.onCancel(ForwardingServerCallListener.java:40) [grpc-api-1.52.1.jar:1.52.1] at io.grpc.Contexts$ContextualizedServerCallListener.onCancel(Contexts.java:96) [grpc-api-1.52.1.jar:1.52.1] at io.grpc.internal.ServerCallImpl$ServerStreamListenerImpl.closedInternal(ServerCallImpl.java:378) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ServerCallImpl$ServerStreamListenerImpl.closed(ServerCallImpl.java:365) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ServerImpl$JumpToApplicationThreadServerStreamListener$1Closed.runInContext(ServerImpl.java:923) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:133) [grpc-core-1.52.1.jar:1.52.1] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) [?:?] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) [?:?]
预期仅客户端-网关连接断开,为何会影响网关-游戏服务器的连接?
更新:Context调整后的新问题
发现问题与Context管理相关,在拦截器中设置Context.fork()后,网关能正常向游戏服务器发消息,但游戏服务器向网关推送流消息时出错,错误日志:
StreamShardingClient.error occurred. CANCELLED: Failed to read message. io.grpc.StatusRuntimeException: CANCELLED: Failed to read message. at io.grpc.Status.asRuntimeException(Status.java:539) ~[grpc-api-1.52.1.jar:1.52.1] at io.grpc.stub.ClientCalls$StreamObserverToCallListenerAdapter.onClose(ClientCalls.java:487) [grpc-stub-1.52.1.jar:1.52.1] at io.grpc.internal.DelayedClientCall$DelayedListener$3.run(DelayedClientCall.java:489) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.DelayedClientCall$DelayedListener.delayOrExecute(DelayedClientCall.java:453) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.DelayedClientCall$DelayedListener.onClose(DelayedClientCall.java:486) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ClientCallImpl.closeObserver(ClientCallImpl.java:576) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ClientCallImpl.access$300(ClientCallImpl.java:70) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1StreamClosed.runInternal(ClientCallImpl.java:757) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1StreamClosed.runInContext(ClientCallImpl.java:736) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37) [grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:133) [grpc-core-1.52.1.jar:1.52.1] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) [?:?] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) [?:?] at java.lang.Thread.run(Thread.java:829) [?:?] Caused by: io.grpc.StatusRuntimeException: CANCELLED: call already cancelled. Use ServerCallStreamObserver.setOnCancelHandler() to disable this exception at io.grpc.Status.asRuntimeException(Status.java:530) ~[grpc-api-1.52.1.jar:1.52.1] at io.grpc.stub.ServerCalls$ServerCallStreamObserverImpl.onNext(ServerCalls.java:366) ~[grpc-stub-1.52.1.jar:1.52.1] at scalasharding.sharding.StreamShardingClient$ResponseObserver.$anonfun$onNext$1(StreamShardingClient.scala:124) ~[classes/:?] at scalasharding.sharding.StreamShardingClient$ResponseObserver.$anonfun$onNext$1$adapted(StreamShardingClient.scala:117) ~[classes/:?] at scala.collection.immutable.Set$Set2.foreach(Set.scala:201) ~[scala-library-2.13.7.jar:?] at scalasharding.sharding.StreamShardingClient$ResponseObserver.onNext(StreamShardingClient.scala:117) ~[classes/:?] at scalasharding.sharding.StreamShardingClient$ResponseObserver.onNext(StreamShardingClient.scala:71) [classes/:?] at io.grpc.stub.ClientCalls$StreamObserverToCallListenerAdapter.onMessage(ClientCalls.java:474) ~[grpc-stub-1.52.1.jar:1.52.1] at io.grpc.internal.DelayedClientCall$DelayedListener.onMessage(DelayedClientCall.java:473) ~[grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1MessagesAvailable.runInternal(ClientCallImpl.java:675) ~[grpc-core-1.52.1.jar:1.52.1] at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1MessagesAvailable.runInContext(ClientCallImpl.java:660) ~[grpc-core-1.52.1.jar:1.52.1] ... 5 more
问题分析与解决
初始问题根因:Context传递导致的取消联动
gRPC的Context会自动在调用链中传递,当客户端断开连接触发网关侧的调用取消时,当前的Context会被标记为取消状态。而网关向游戏服务器发起的客户端流调用默认继承了这个已取消的Context,导致游戏服务器的连接也被联动取消。
解决方案步骤
隔离网关-游戏服务器流的Context
在创建网关到游戏服务器的双向流时,不要继承客户端连接的Context,而是创建一个独立的Context:// 替换原有的streamConnection调用部分 val independentContext = Context.current().fork().withCancellation() val streamConnection = independentContext.run(() => getOrCreateRoomConnection(`Context-RoomId-Key`))这里通过
fork()创建独立上下文,并确保它不受客户端连接上下文取消的影响。处理客户端连接断开的清理逻辑
当客户端断开时,只需要清理该客户端对应的网关侧响应观察者,不要影响共享的网关-游戏服务器流:override def onCancel(): Unit = { // 仅移除当前客户端的响应观察者,不关闭网关-游戏服务器的流 removeClientObserver(userId, responseObserver) // 禁用默认的取消传播 val serverStreamObserver = responseObserver.asInstanceOf[ServerCallStreamObserver[_]] serverStreamObserver.setOnCancelHandler(() => { // 自定义取消逻辑,不传播到下游 }) }修复Context fork后的反向流问题
更新后的错误提示call already cancelled,是因为游戏服务器推送消息时,尝试使用已取消的客户端上下文发送响应。需要确保游戏服务器的响应流使用独立的上下文,或者在转发响应时切换到客户端的上下文但忽略取消状态:// 在游戏服务器响应的处理逻辑中,使用独立上下文发送到客户端 val clientResponseContext = Context.current().fork() clientResponseContext.run(() => { responseObserver.onNext(gameServerResponse) })或者,针对客户端的
ServerCallStreamObserver设置取消处理器,避免在客户端断开后发送响应时触发错误:val serverObserver = responseObserver.asInstanceOf[ServerCallStreamObserver[StateSyncFrameDemo]] serverObserver.setOnCancelHandler(() => { // 标记该客户端已断开,后续不再发送消息 markClientDisconnected(userId) })
核心总结
- gRPC的
Context取消状态会自动传播,必须隔离共享流与客户端流的上下文 - 使用
Context.fork()创建独立上下文,避免联动取消 - 针对
ServerCallStreamObserver设置自定义取消处理器,禁用默认的取消传播逻辑
内容的提问来源于stack exchange,提问作者LoranceChen

