【问题标题】:gRPC: Rendezvous terminated with (StatusCode.INTERNAL, Received RST_STREAM with error code 2)gRPC:集合以(StatusCode.INTERNAL,收到错误代码 2 的 RST_STREAM)终止
【发布时间】:2018-06-18 20:17:12
【问题描述】:

我正在 python 中实现 gRPC 客户端和服务器。服务器成功从客户端接收数据,但客户端收到“RST_STREAM with error code 2”。

这实际上是什么意思,我该如何解决?

这是我的原型文件:

service MyApi {
    rpc SelectModelForDataset (Dataset) returns (SelectedModel) {
    }
}
message Dataset {
    // ...
}
message SelectedModel {
    // ...
}

我的服务实现如下所示:

class MyApiServicer(my_api_pb2_grpc.MyApiServicer):
def SelectModelForDataset(self, request, context):
    print("Processing started.")
    selectedModel = ModelSelectionModule.run(request, context)  
    print("Processing Completed.")
    return selectedModel

我用这段代码启动服务器:

import grpc
from concurrent import futures
#...
server = grpc.server(futures.ThreadPoolExecutor(max_workers=100))
my_api_pb2_grpc.add_MyApiServicer_to_server(MyApiServicer(), server)
server.add_insecure_port('[::]:50051')
server.start()

我的客户是这样的:

channel = grpc.insecure_channel(target='localhost:50051')
stub = my_api_pb2_grpc.MyApiStub(channel)
dataset = my_api_pb2.Dataset() 
# fill the object ...
model = stub.SelectModelForDataset(dataset)  # call server

客户端调用后,服务器开始处理直到完成(大约需要一分钟),但客户端立即返回并出现以下错误:

Traceback (most recent call last):                                                                   
File "Client.py", line 32, in <module>                                                               
    run()                                                                                            
File "Client.py", line 26, in run                                                                    
    model = stub.SelectModelForDataset(dataset)  # call server                                       
File "/usr/local/lib/python3.5/dist-packages/grpc/_channel.py", line 484, in __call__
    return _end_unary_response_blocking(state, call, False, deadline)                                
File "/usr/local/lib/python3.5/dist-packages/grpc/_channel.py", line 434, in _end_unary_response_blocking                                                                                               
    raise _Rendezvous(state, None, None, deadline)                                                 
grpc._channel._Rendezvous: <_Rendezvous of RPC that terminated with (StatusCode.INTERNAL, Received RST_STREAM with error code 2)>

如果我异步执行请求并等待未来,

model_future = stub.SelectModelForDataset.future(dataset)  # call server
model = model_future.result()

客户端一直等到完成,但之后仍然返回错误:

Traceback (most recent call last):                                                                   
File "AsyncClient.py", line 35, in <module>                                                          
    run()                                                                                            
File "AsyncClient.py", line 29, in run                                                               
    model = model_future.result()                                                                    
File "/usr/local/lib/python3.5/dist-packages/grpc/_channel.py", line 276, in result                  
    raise self                                                                                     
grpc._channel._Rendezvous: <_Rendezvous of RPC that terminated with (StatusCode.INTERNAL, Received RST_STREAM with error code 2)>

UPD:启用跟踪 GRPC_TRACE=all 后,我发现了以下内容:

客户端,请求后立即:

E0109 17:59:42.248727600    1981 channel_connectivity.cc:126] watch_completion_error: {"created":"@1515520782.248638500","description":"GOAWAY received","file":"src/core/ext/transport/chttp2/transport/chttp2_transport.cc","file_line":1137,"http2_error":0,"raw_bytes":"Server shutdown"}            
E0109 17:59:42.451048100    1979 channel_connectivity.cc:126] watch_completion_error: "Cancelled"  
E0109 17:59:42.451160000    1979 completion_queue.cc:659]    Operation failed: tag=0x7f6e5cd1caf8, error={"created":"@1515520782.451034300","description":"Timed out waiting for connection state change","file":"src/core/ext/filters/client_channel/channel_connectivity.cc","file_line":133}
...(last two messages keep repeating 5 times every second)

服务器:

E0109 17:59:42.248201000    1985 completion_queue.cc:659]    Operation failed: tag=0x7f3f74febee8, error={"created":"@1515520782.248170000","description":"Server Shutdown","file":"src/core/lib/surface/server.cc","file_line":1249}                                                                    
E0109 17:59:42.248541100    1975 tcp_server_posix.cc:231]    Failed accept4: Invalid argument                                                                             
E0109 17:59:47.362868700    1994 completion_queue.cc:659]    Operation failed: tag=0x7f3f74febee8, error={"created":"@1515520787.362853500","description":"Server Shutdown","file":"src/core/lib/surface/server.cc","file_line":1249}                                                                                                                                             
E0109 17:59:52.430612500    2000 completion_queue.cc:659]    Operation failed: tag=0x7f3f74febee8, error={"created":"@1515520792.430598800","description":"Server Shutdown","file":"src/core/lib/surface/server.cc","file_line":1249}
... (last message kept repeating every few seconds)                                                             

UPD2:

我的Server.py文件的全部内容:

import ModelSelectionModule
import my_api_pb2_grpc
import my_api_pb2
import grpc
from concurrent import futures
import time

class MyApiServicer(my_api_pb2_grpc.MyApiServicer):
    def SelectModelForDataset(self, request, context):
        print("Processing started.")
        selectedModel = ModelSelectionModule.run(request, context)
        print("Processing Completed.")
        return selectedModel


# TODO(shalamov): what is the best way to run a python server?
def serve():
    server = grpc.server(futures.ThreadPoolExecutor(max_workers=100))
    my_api_pb2_grpc.add_MyApiServicer_to_server(MyApiServicer(), server)
    server.add_insecure_port('[::]:50051')
    server.start()

    print("gRPC server started\n")
    try:
        while True:
            time.sleep(24 * 60 * 60)  # run for 24h
    except KeyboardInterrupt:
        server.stop(0)


if __name__ == '__main__':
    serve()

UPD3: 似乎,ModelSelectionModule.run 导致了问题。我试图将它隔离到一个单独的线程中,但它没有帮助。 selectedModel 最终会被计算出来,但是那个时候客户端已经不在了。 如何防止此调用与 grpc 混淆?

pool = ThreadPool(processes=1)
async_result = pool.apply_async(ModelSelectionModule.run(request, context))
selectedModel = async_result.get()

调用相当复杂,它产生并加入许多线程,调用不同的库,如scikit-learnsmac 等。全部发在这里就太过分了。

在调试时,我发现在客户端请求后,服务器保持 2 个连接打开(fd 3fd 8)。如果我手动关闭fd 8 或向其写入一些字节,我在客户端中看到的错误将变为Stream removed(而不是Received RST_STREAM with error code 2)。似乎,套接字(fd 8)不知何故被子进程损坏了。 这怎么可能?如何保护套接字不被子进程访问?

【问题讨论】:

  • 你的 server.start() 之后有 sleep() 函数吗? server.start() 在后台线程中启动服务器,如果程序在启动服务器后退出,该线程将进行垃圾收集。
  • 当然,server.start() 之后我有睡眠功能。服务器正在运行并处理每个请求,但客户端在请求后立即收到这个神秘的 RST_STREAM。
  • 你能发布你的 server.start() 文件的全部内容吗?这看起来很像一个服务器被过早地销毁。
  • 服务器设置看起来正确。如果您将 ModelSelectionModule.run(request, context) 函数替换为睡眠(在相似的时间量内),并返回一个虚拟的 ModelSelectionModule,这行得通吗?可能是 ModelSelectionModule.run fork() 工作并最终弄乱了文件描述符。
  • 在处理程序中使用 fork() 最安全的方法可能是 fork-then-exec。如果这不可能,请确保在子进程返回 grpc 函数处理程序之前退出()子进程可能会有所帮助。

标签: python python-3.x protocol-buffers grpc


【解决方案1】:

这是在进程处理程序中使用 fork() 的结果。 gRPC Python 不支持此用例。

【讨论】:

  • 听到这个我很惊讶。为什么 fork() 会断开连接?有办法吗?我的代码去年 9 月运行良好,现在与 RST_STREAM 中断,这对我的项目来说非常有问题。
【解决方案2】:

我遇到了这个问题,刚刚解决了,你用的是with_call()方法吗?

错误代码:

response = stub.SayHello.with_call(request=request, metadata=metadata)

响应是一个元组。

成功代码:不要使用 with_call()

response = stub.SayHello(request=request, metadata=metadata)

响应是一个响应对象。

【讨论】:

    猜你喜欢
    • 2021-10-30
    • 2022-09-28
    • 2021-10-26
    • 1970-01-01
    • 2021-12-04
    • 2023-02-17
    • 2012-06-29
    • 2019-05-23
    • 1970-01-01
    相关资源
    最近更新 更多