Python gRPC客户端双向流中,等价于Go的stream.Send()的方法是什么?
Python gRPC双向流客户端实现同流回发请求的方法
在Python的gRPC中,双向流的客户端实现可以通过调用RouteChat返回的流对象的send()方法,在同一条流里主动向服务器发送消息,无需依赖一次性生成器传参。
改造思路
原代码通过生成器一次性发送所有请求后被动接收响应,要实现动态回发需调整为:
- 调用
stub.RouteChat()时传入空生成器(如iter([])),获取双向流对象。 - 单独启动线程监听服务器响应,避免阻塞发送操作。
- 主线程通过流对象的
send()方法,随时发送新的请求消息。
改造后的示例代码
import threading import time from your_proto_module import RouteNote, make_route_note # 替换为你的proto实际模块 def receive_responses(stream): """监听服务器响应的线程函数""" try: for response in stream: print(f"Received message '{response.message}' at ({response.location.latitude}, {response.location.longitude})") except Exception as e: print(f"接收响应出错: {e}") def guide_route_chat(stub): # 初始化双向流,传入空生成器 stream = stub.RouteChat(iter([])) # 启动线程异步监听响应 threading.Thread(target=receive_responses, args=(stream,), daemon=True).start() # 发送初始消息 initial_messages = [ make_route_note("First message", 0, 0), make_route_note("Second message", 0, 1), ] for msg in initial_messages: print(f"Sending '{msg.message}' at ({msg.location.latitude}, {msg.location.longitude})") stream.send(msg) time.sleep(1) # 模拟业务场景下的动态回发(比如收到响应后、用户输入触发等) time.sleep(3) new_msg = make_route_note("Dynamic reply message", 1, 1) print(f"Sending dynamic reply: '{new_msg.message}' at ({new_msg.location.latitude}, {new_msg.location.longitude})") stream.send(new_msg) # 根据业务需求决定何时关闭流 time.sleep(2) stream.close()
关键说明
- 双向流对象为
grpc._channel._MultiThreadedRendezvous类型,自带send()方法用于发送消息,close()方法用于关闭客户端发送流。 - 必须用独立线程处理响应接收,因为
for response in stream是阻塞操作,主线程执行会无法同时发送消息。 - 发送消息前需确保流未关闭,否则会触发异常。
内容的提问来源于stack exchange,提问作者chinmay
相关产品推荐
相关产品推荐

