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

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两个节点转发相同的请求流,并且将两个节点的响应实时返回给本地客户端。

当前写法存在两个核心问题:

  1. gRPC的请求迭代器request_iterator只能被遍历一次,传给node1后,node2无法获取任何请求数据。
  2. 串行调用两个节点的RouteChat接口,会导致node2的响应必须等node1的流结束后才能返回,无法实现实时响应。

解决方案

实现思路

  1. 缓存请求数据:先将客户端发来的所有请求消息缓存(如果是无限流,改用队列+线程实时转发),确保两个节点都能拿到完整请求流。
  2. 并发调用多节点:使用线程池同时向两个节点发起RouteChat调用,并行处理请求。
  3. 合并响应流:实时收集两个节点的响应,一旦有响应就立即返回给客户端,保证响应实时性。

修改后的服务端代码(有限流场景)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 08:15:43