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
相关产品推荐
相关产品推荐

