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

Python gRPC双向流服务:带客户端交互的解决方案生成实现问题

Python gRPC双向流服务实现困境与解决方案

问题描述

我正在用Python和gRPC构建双向调用服务,客户端能通过该服务请求服务器的计算解决方案流,而且同一连接要支持多种消息交互。核心逻辑是:客户端发送NextMsg请求解决方案时,服务端得先向客户端请求特定信息,完成计算后再返回结果。

我尝试用抽象类通用化解方案生成流程,但Python gRPC的双向流方法用的是迭代器,没法像Kotlin那样传递响应观察者引用,再加上项目要求混合异步同步编程,现在卡壳了,求解决办法。

.proto文件骨架

service GenericPrimitiveService {
  rpc callPrimitive(stream ClientMsg) returns (stream ServerMsg) {}
}

message ClientMsg {
    oneof msg {
        string request = 1;
        string next = 2;
        // 新增客户端回复服务端请求的字段
        string info = 3;
    }
}

message ServerMsg {
    oneof msg {
        string solution = 1;
        // 新增服务端向客户端请求信息的字段
        string ask_for_info = 2;
    }
}

预期代码逻辑

def callPrimitive(self, request_iterator, context):
        stream = None 
        for msg in request_iterator:
            if(stream == None and msg.request.IsInitialized()):
                stream = primitive(msg.request)
            elif(msg.next.IsInitialized()):
                yield next(stream)
            else:
                handleReceivedInformation(msg)

def primitive(request):
    while True:
        # 向双向流的生成器发送请求信息给客户端
        response = await sendRequestToClient()
        solution = elaborate(response)
        yield(solution)

@abstractmethod
def elaborate(info):
    pass

解决方案

核心思路

Python gRPC双向流的本质是请求迭代器(服务端接收客户端消息)与响应生成器(服务端发送消息给客户端)的分离,无法直接传递观察者引用,因此需要通过中间通信层(队列+事件)实现服务端主动向客户端请求信息的逻辑,同时用异步IO协调同步与异步代码。

具体实现

1. 抽象基类封装通用逻辑

定义抽象基类,封装向客户端请求信息、等待回复的通用流程,子类仅需实现具体计算逻辑:

import asyncio
from collections import deque
from abc import ABC, abstractmethod

class BasePrimitive(ABC):
    def __init__(self, request):
        self.request = request
        # 存储客户端回复的信息
        self.client_response_queue = deque()
        # 标记客户端是否已回复
        self.response_received = asyncio.Event()

    async def _request_client_input(self):
        # 生成服务端向客户端请求信息的消息
        return ServerMsg(ask_for_info=f"请提供{self.request}相关的补充信息")

    @abstractmethod
    def elaborate(self, info):
        # 子类实现具体计算逻辑
        pass

    async def generate_solutions(self):
        while True:
            # 第一步:向客户端发送请求信息
            request_msg = await self._request_client_input()
            yield request_msg
            
            # 第二步:等待客户端回复
            self.response_received.clear()
            await self.response_received.wait()
            
            # 第三步:获取回复并计算解决方案
            client_info = self.client_response_queue.popleft()
            solution = self.elaborate(client_info)
            yield ServerMsg(solution=solution)

2. 实现具体Primitive子类

class ConcretePrimitive(BasePrimitive):
    def elaborate(self, info):
        # 示例计算逻辑:根据客户端提供的信息生成解决方案
        return f"基于请求[{self.request}]和补充信息[{info}]生成的解决方案"

3. gRPC服务实现

在gRPC服务方法中,协调请求迭代器、响应生成器与Primitive实例的通信:

import grpc
from concurrent import futures
# 导入生成的proto代码(替换为你的实际模块名)
from your_proto_pb2 import ClientMsg, ServerMsg
from your_proto_pb2_grpc import GenericPrimitiveServiceServicer, add_GenericPrimitiveServiceServicer_to_server

class GenericPrimitiveService(GenericPrimitiveServiceServicer):
    def callPrimitive(self, request_iterator, context):
        primitive = None
        solution_gen = None
        loop = asyncio.get_event_loop()
        pending_messages = []

        for msg in request_iterator:
            # 初始化Primitive实例
            if msg.HasField("request") and not primitive:
                primitive = ConcretePrimitive(msg.request)
                solution_gen = primitive.generate_solutions()
                # 预取第一个要发送给客户端的请求消息
                pending_messages.append(loop.run_until_complete(asyncio.ensure_future(solution_gen.__anext__())))
            
            # 处理客户端的Next请求,发送解决方案或请求信息
            elif msg.HasField("next") and primitive:
                if pending_messages:
                    yield pending_messages.pop(0)
                else:
                    try:
                        next_msg = loop.run_until_complete(asyncio.ensure_future(solution_gen.__anext__()))
                        yield next_msg
                    except StopAsyncIteration:
                        break
            
            # 处理客户端回复的信息,通知Primitive实例
            elif msg.HasField("info") and primitive:
                primitive.client_response_queue.append(msg.info)
                primitive.response_received.set()

# 启动gRPC服务
def serve():
    server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
    add_GenericPrimitiveServiceServicer_to_server(GenericPrimitiveService(), server)
    server.add_insecure_port("[::]:50051")
    server.start()
    server.wait_for_termination()

if __name__ == "__main__":
    serve()

关键说明

  • 队列与事件机制:用deque存储客户端回复,asyncio.Event通知生成器继续执行,解决了同步迭代器与异步生成器之间的通信问题;
  • 异步生成器:将解决方案生成逻辑封装为异步生成器,既支持向客户端请求的异步操作,又能通过yield向gRPC响应流发送消息;
  • 抽象类复用:BasePrimitive封装了通用的请求-等待-计算流程,子类只需实现elaborate即可扩展不同的计算逻辑;
  • 混合编程协调:在同步的gRPC方法中,通过asyncio事件循环运行异步任务,实现同步与异步代码的无缝衔接。

内容的提问来源于stack exchange,提问作者L0L_Osi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 01:47:16