如何使用Python gRPC处理流式消息

evi*_*obu 2 python stream grpc

我正在关注此Route_Guide示例

有问题的示例将触发并读取消息,而无需回复特定的消息。后者是我想要实现的目标。

这是我到目前为止的内容:

import grpc
...

channel = grpc.insecure_channel(conn_str)
try:
    grpc.channel_ready_future(channel).result(timeout=5)
except grpc.FutureTimeoutError:
    sys.exit('Error connecting to server')
else:
    stub = MyService_pb2_grpc.MyServiceStub(channel)
    print('Connected to gRPC server.')
    this_is_just_read_maybe(stub)


def this_is_just_read_maybe(stub):
    responses = stub.MyEventStream(stream())
    for response in responses:
        print(f'Received message: {response}')
        if response.something:
            # okay, now what? how do i send a message here?

def stream():
    yield my_start_stream_msg
    # this is fine, i receive this server-side
    # but i can't check for incoming messages here
Run Code Online (Sandbox Code Playgroud)

我似乎在存根上没有read()write(),所有内容似乎都是通过迭代器实现的。

我该如何发送消息this_is_just_read_maybe(stub)?那是正确的方法吗?

我的Proto是双向流:

service MyService {
  rpc MyEventStream (stream StreamingMessage) returns (stream StreamingMessage) {}
}
Run Code Online (Sandbox Code Playgroud)

小智 6

您尝试做的事情完全有可能实现,并且可能涉及编写自己的请求迭代器对象,该对象可以在到达时得到响应,而不是使用简单的生成器作为请求迭代器。也许像

class MySmarterRequestIterator(object):

    def __init__(self):
        self._lock = threading.Lock()
        self._responses_so_far = []

    def __iter__(self):
        return self

    def _next(self):
        # some logic that depends upon what responses have been seen
        # before returning the next request message
        return <your message value>

    def __next__(self):  # Python 3
        return self._next()

    def next(self):  # Python 2
        return self._next()

    def add_response(self, response):
        with self._lock:
            self._responses.append(response)
Run Code Online (Sandbox Code Playgroud)

然后你使用像

my_smarter_request_iterator = MySmarterRequestIterator()
responses = stub.MyEventStream(my_smarter_request_iterator)
for response in responses:
    my_smarter_request_iterator.add_response(response)
Run Code Online (Sandbox Code Playgroud)

。您的_next实现中可能会存在锁定和阻塞,以处理gRPC Python的情况,即询问您的对象该对象要发送的下一个请求以及您的响应(实际上)是“等待,请稍等,我不知道我的请求是什么希望发送,直到我看到下一个响应的结果。”

  • 问题是,当使用所有这些锁为许多客户端提供双向连接时,此解决方案的可伸缩性如何? (2认同)

dco*_*les 6

除了编写自定义迭代器,您还可以使用阻塞队列为客户端存根实现类似发送和接收的行为:

import queue
...

send_queue = queue.SimpleQueue()  # or Queue if using Python before 3.7
my_event_stream = stub.MyEventStream(iter(send_queue.get, None))

# send
send_queue.push(StreamingMessage())

# receive
response = next(my_event_stream)  # type: StreamingMessage
Run Code Online (Sandbox Code Playgroud)

这利用了 的标记形式iter,它将常规函数转换为迭代器,该迭代器在达到标记值时停止(在本例中为None)。