grpc.aio中StreamStreamClientInterceptor的intercept_stream_stream未被调用
gRPC Python异步客户端StreamStream拦截器未触发问题解决
问题现象
同步gRPC客户端中,UnaryUnaryClientInterceptor和StreamStreamClientInterceptor的拦截方法均能正常触发;但切换到异步(grpc.aio)版本后,仅UnaryUnaryClientInterceptor的intercept_unary_unary被调用,StreamStreamClientInterceptor的intercept_stream_stream完全不执行。
问题原因
异步版本中intercept_stream_stream方法的实现存在错误:对continuation的返回值执行了await操作。实际上,StreamStream类型的异步调用中,continuation返回的是异步迭代器对象(而非可await的Future),await该对象会导致拦截逻辑提前终止,无法传递到后续的异步遍历流程,最终表现为拦截器未被触发。
解决方案
修改异步拦截器的intercept_stream_stream方法,移除await,直接返回continuation的调用结果。
修改后的异步拦截器代码
import asyncio import grpc from grpc.aio import ClientCallDetails, AioRpcError import geyser_pb2 import geyser_pb2_grpc class WithHeaders( grpc.aio.UnaryUnaryClientInterceptor, grpc.aio.StreamStreamClientInterceptor, ): def __init__(self, *headers: tuple[str, str]): self.headers = headers def _insert_headers(self, new_metadata, client_call_details) -> ClientCallDetails: metadata = [] if client_call_details.metadata is not None: metadata = list(client_call_details.metadata) metadata.extend(new_metadata) return ClientCallDetails( method=client_call_details.method, timeout=client_call_details.timeout, metadata=metadata, credentials=client_call_details.credentials, wait_for_ready=client_call_details.wait_for_ready, ) async def intercept_unary_unary(self, continuation, client_call_details, request): print("intercept unary_unary") new_client_call_details = self._insert_headers( self.headers, client_call_details ) return await continuation(new_client_call_details, request) async def intercept_stream_stream( self, continuation, client_call_details, request_iterator ): print("intercept stream_stream") new_client_call_details = self._insert_headers( self.headers, client_call_details ) # 移除await,直接返回continuation的结果 return continuation(new_client_call_details, request_iterator) # 其余GeyserClient和main代码保持不变
关键说明
- 对于Unary-Unary类型的异步调用,
continuation返回的是Future对象,必须通过await获取最终响应; - 对于Stream-Stream类型的异步调用,
continuation返回的是异步调用实例(本身是异步迭代器),直接返回即可,由后续的async for遍历触发实际调用流程。
内容的提问来源于stack exchange,提问作者Martim Martins
相关产品推荐
相关产品推荐

