Python gRPC流式服务端如何生成双客户端连接同服务端并实时返回响应
gRPC双向流转发:多节点同步接收与响应实时返回
现有代码
1. Proto定义
service RouteGuide { rpc RouteChat(stream RouteNote) returns (stream RouteNote) {} } message RouteNote { // The location from which the message is sent. Point location = 1; // The message to be sent. string message = 2; }
2. 本地客户端代码
def generate_messages(): messages = [ make_route_note("First message", 0, 0), make_route_note("Second message", 0, 1), make_route_note("Third message", 1, 0), make_route_note("Fourth message", 0, 0), make_route_note("Fifth message", 1, 0), ] for msg in messages: print("Sending %s at %s" % (msg.message, msg.location)) yield msg def guide_route_chat(stub): responses = stub.RouteChat(generate_messages()) for response in responses: print("Received message %s at %s" % (response.message, response.location)) def run(): with grpc.insecure_channel('localhost:50051') as channel: stub = route_guide_pb2_grpc.RouteGuideStub(channel) print("-------------- RouteChat --------------") guide_route_chat(stub) if __name__ == '__main__': run()
3. 原服务端代码
class RouteGuideServicer(route_guide_pb2_grpc.RouteGuideServicer): def RouteChat(self, request_iterator, context): yield from request_iterator def serve(): server = grpc.server(futures.ThreadPoolExecutor(max_workers=10)) route_guide_pb2_grpc.add_RouteGuideServicer_to_server(RouteGuideServicer(), server) server.add_insecure_port('[::]:50051') server.start() server.wait_for_termination() if __name__ == '__main__': serve()
需求与当前问题
需要在服务端的RouteGuideServicer.RouteChat方法中,同时向node1:50051和node2:50051两个节点转发相同的请求流,并且将两个节点的响应实时返回给本地客户端。
当前写法存在两个核心问题:
- gRPC的请求迭代器
request_iterator只能被遍历一次,传给node1后,node2无法获取任何请求数据。 - 串行调用两个节点的
RouteChat接口,会导致node2的响应必须等node1的流结束后才能返回,无法实现实时响应。
解决方案
实现思路
- 缓存请求数据:先将客户端发来的所有请求消息缓存(如果是无限流,改用队列+线程实时转发),确保两个节点都能拿到完整请求流。
- 并发调用多节点:使用线程池同时向两个节点发起
RouteChat调用,并行处理请求。 - 合并响应流:实时收集两个节点的响应,一旦有响应就立即返回给客户端,保证响应实时性。
修改后的服务端代码(有限流场景)
import grpc import route_guide_pb2 import route_guide_pb2_grpc from concurrent import futures import threading from queue import Queue class RouteGuideServicer(route_guide_pb2_grpc.RouteGuideServicer): def RouteChat(self, request_iterator, context): # 缓存所有请求消息 requests = list(request_iterator) # 用队列收集两个节点的响应 response_queue = Queue() def call_remote_node(node_addr): """向单个远程节点发送请求并收集响应""" try: with grpc.insecure_channel(node_addr) as channel: stub = route_guide_pb2_grpc.RouteGuideStub(channel) # 用缓存的请求列表生成迭代器传给远程节点 for resp in stub.RouteChat(iter(requests)): response_queue.put(resp) except Exception as e: print(f"Error calling {node_addr}: {e}") # 启动线程同时调用两个节点 nodes = ['node1:50051', 'node2:50051'] threads = [] for node in nodes: t = threading.Thread(target=call_remote_node, args=(node,)) t.start() threads.append(t) # 实时从队列取出响应并返回给客户端 active_threads = len(threads) while active_threads > 0 or not response_queue.empty(): try: resp = response_queue.get(timeout=0.1) yield resp except: # 更新活跃线程数 active_threads = sum(1 for t in threads if t.is_alive()) # 等待所有线程结束 for t in threads: t.join() def serve(): server = grpc.server(futures.ThreadPoolExecutor(max_workers=10)) route_guide_pb2_grpc.add_RouteGuideServicer_to_server(RouteGuideServicer(), server) server.add_insecure_port('[::]:50051') server.start() server.wait_for_termination() if __name__ == '__main__': serve()
无限流场景优化(可选)
如果客户端发送的是持续的无限流,不能一次性缓存所有请求,可改为实时转发模式:
class RouteGuideServicer(route_guide_pb2_grpc.RouteGuideServicer): def RouteChat(self, request_iterator, context): request_queue = Queue() response_queue = Queue() done_event = threading.Event() def feed_requests(): """将客户端请求实时放入队列""" try: for req in request_iterator: request_queue.put(req) finally: # 请求流结束后标记 request_queue.put(None) done_event.set() def call_remote_node(node_addr): """从队列取请求并发给远程节点,收集响应""" try: with grpc.insecure_channel(node_addr) as channel: stub = route_guide_pb2_grpc.RouteGuideStub(channel) def request_generator(): while True: req = request_queue.get() if req is None: break yield req for resp in stub.RouteChat(request_generator()): response_queue.put(resp) except Exception as e: print(f"Error calling {node_addr}: {e}") # 启动请求转发线程 threading.Thread(target=feed_requests, daemon=True).start() # 启动两个节点调用线程 nodes = ['node1:50051', 'node2:50051'] for node in nodes: threading.Thread(target=call_remote_node, args=(node,), daemon=True).start() # 实时返回响应 while not done_event.is_set() or not response_queue.empty(): try: resp = response_queue.get(timeout=0.1) yield resp except: pass
关键说明
- 请求复用:通过列表或队列解决迭代器只能遍历一次的问题,确保两个节点都能获取完整请求。
- 并发处理:使用线程并行调用两个节点,避免串行等待导致的响应延迟。
- 响应实时性:通过队列收集响应,一旦有响应就立即返回给客户端,保证实时性。
内容的提问来源于stack exchange,提问作者iskylite
相关产品推荐
相关产品推荐

