You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java gRPC双向流异常:客户端取消致网关-游戏服务连接中断

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,导致游戏服务器的连接也被联动取消。

解决方案步骤

  1. 隔离网关-游戏服务器流的Context
    在创建网关到游戏服务器的双向流时,不要继承客户端连接的Context,而是创建一个独立的Context:

    // 替换原有的streamConnection调用部分
    val independentContext = Context.current().fork().withCancellation()
    val streamConnection = independentContext.run(() => getOrCreateRoomConnection(`Context-RoomId-Key`))
    

    这里通过fork()创建独立上下文,并确保它不受客户端连接上下文取消的影响。

  2. 处理客户端连接断开的清理逻辑
    当客户端断开时,只需要清理该客户端对应的网关侧响应观察者,不要影响共享的网关-游戏服务器流:

    override def onCancel(): Unit = {
      // 仅移除当前客户端的响应观察者,不关闭网关-游戏服务器的流
      removeClientObserver(userId, responseObserver)
      // 禁用默认的取消传播
      val serverStreamObserver = responseObserver.asInstanceOf[ServerCallStreamObserver[_]]
      serverStreamObserver.setOnCancelHandler(() => {
        // 自定义取消逻辑,不传播到下游
      })
    }
    
  3. 修复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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 02:24:58