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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 15:50:00