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

基于grpc.aio的Python服务端请求耗时日志拦截器实现求助

gRPC AIO服务端请求耗时统计拦截器实现

针对你遇到的问题,下面是完整的异步gRPC服务端日志拦截器实现,能够正确统计请求处理耗时:

完整的grpc.py代码

import asyncio
from typing import Callable, Awaitable

import grpc

from pkg.py.gen import service_pb2
from pkg.py.gen import service_pb2_grpc
from ..logger import get_logger
from ..services import Service


class SampleService(service_pb2_grpc.SampleServiceServicer):
    def __init__(self) -> None:
        super().__init__()

    async def Hello(self, request: service_pb2.HelloRequest, _) -> service_pb2.HelloResponse:
        await asyncio.sleep(1)
        return service_pb2.HelloResponse(Greeting=Service.hello(request.Name))


class LoggingServerInterceptor(grpc.aio.ServerInterceptor):
    async def intercept_service(
            self,
            continuation: Callable[
                [grpc.HandlerCallDetails], Awaitable[grpc.RpcMethodHandler]
            ],
            handler_call_details: grpc.HandlerCallDetails,
    ) -> grpc.RpcMethodHandler:
        logger = get_logger()
        # 获取原始的RPC方法处理器
        handler = await continuation(handler_call_details)

        # 仅处理存在invoke方法的Unary-Unary类型RPC
        if handler and handler.invoke is not None:
            original_invoke = handler.invoke

            async def wrapped_invoke(request_or_iterator, context):
                # 记录请求开始时间
                start_time = asyncio.get_event_loop().time()
                try:
                    # 等待原始请求处理完成
                    result = await original_invoke(request_or_iterator, context)
                    return result
                finally:
                    # 计算耗时并记录日志
                    duration = asyncio.get_event_loop().time() - start_time
                    logger.info(
                        f"gRPC request processed | method={handler_call_details.method} | duration={round(duration*1000, 2)}ms | status={context.code().name}",
                        extra={
                            "method": handler_call_details.method,
                            "duration_ms": round(duration * 1000, 2),
                            "status_code": context.code().value[0],
                            "status_name": context.code().name
                        }
                    )

            # 替换原始处理器的invoke方法为包装后的方法
            handler = handler._replace(invoke=wrapped_invoke)
        
        return handler

关键实现说明

  1. 获取原始处理器:通过await continuation(handler_call_details)获取后续拦截器或实际服务的处理器,必须使用await确保拿到有效处理器。
  2. 包装invoke方法:针对Unary-Unary类型的RPC,我们包装其invoke方法——这是实际处理请求的入口。
  3. 正确计时:在调用原始invoke前记录开始时间,await原始方法确保等待请求处理完成后再计算耗时,解决你之前遇到的“continuation未等待请求完成”问题。
  4. 日志上下文:记录请求方法名、耗时(毫秒)、状态码等关键信息,方便后续排查问题。

其他RPC类型支持(可选)

如果你的服务包含流式RPC(Unary-Stream、Stream-Unary、Stream-Stream),需要额外包装对应的stream方法,示例如下:

# 在LoggingServerInterceptor中,添加stream方法的包装逻辑
if handler and handler.stream is not None:
    original_stream = handler.stream

    async def wrapped_stream(request_or_iterator, context):
        start_time = asyncio.get_event_loop().time()
        try:
            async for response in original_stream(request_or_iterator, context):
                yield response
        finally:
            duration = asyncio.get_event_loop().time() - start_time
            logger.info(
                f"gRPC stream processed | method={handler_call_details.method} | duration={round(duration*1000, 2)}ms | status={context.code().name}"
            )

    handler = handler._replace(stream=wrapped_stream)

验证效果

启动服务后调用Hello接口,日志会输出类似如下内容:

INFO: ...: gRPC request processed | method=/auth.SampleService/Hello | duration=1001.23ms | status=OK

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 01:55:27