跳转至

Python 同步与异步 gRPC 开发指南

Python 在数据分析、AI 大模型推理(如 PyTorch / vLLM / Triton)与自动化任务中占据主导地位。

在 Python 体系中,gRPC 提供了传统的多线程同步模式(Thread-Pool based)与现代的异步原生协程模式(grpc.aio。在高并发 I/O 密集型场景下,推荐全面拥抱 grpc.aio


1. 环境准备与依赖安装

在虚拟环境(venv)中安装 Python gRPC 工具链:

pip install grpcio grpcio-tools grpcio-reflection grpcio-health-checking

代码生成脚本:

python -m grpc_tools.protoc \
  -I. \
  --python_out=. \
  --grpc_python_out=. \
  api/v1/order.proto


2. 现代异步服务端实战 (grpc.aio)

基于 Python 3.8+ asyncio 原生协程的高性能服务端实现:

import asyncio
import logging
import signal
import grpc
from grpc_reflection.v1alpha import reflection

# 导入编译生成的桩文件
import order_pb2
import order_pb2_grpc

class AsyncOrderService(order_pb2_grpc.OrderServiceServicer):
    """基于 asyncio 的订单服务实现"""

    async def CreateOrder(self, request, context):
        if not request.user_id:
            await context.abort(grpc.StatusCode.INVALID_ARGUMENT, "user_id 不能为空")

        order_id = f"py_ord_{int(asyncio.get_event_loop().time() * 1000)}"
        logging.info("成功创建订单: %s 为用户: %s", order_id, request.user_id)

        return order_pb2.CreateOrderResponse(
            order_id=order_id,
            status="CREATED"
        )

    async def TrackOrder(self, request, context):
        steps = ["订单已接收", "仓库打包装箱", "顺丰冷链装车", "派送中", "已送达"]
        for step in steps:
            # 检查客户端是否已取消流
            if context.cancelled():
                logging.warning("客户端已断开订阅,提前终止推流")
                return

            yield order_pb2.TrackOrderResponse(
                status=step,
                message=f"订单 {request.order_id} 进度: {step}",
                timestamp=int(asyncio.get_event_loop().time())
            )
            # 异步非阻塞休眠
            await asyncio.sleep(0.8)

async def serve():
    # 创建异步 gRPC Server (无需占用庞大线程池,单进程承载万级并发 Stream)
    server = grpc.aio.server()
    order_pb2_grpc.add_OrderServiceServicer_to_server(AsyncOrderService(), server)

    # 启用反射服务便于调试
    service_names = (
        order_pb2.DESCRIPTOR.services_by_name['OrderService'].full_name,
        reflection.SERVICE_NAME,
    )
    reflection.enable_server_reflection(service_names, server)

    listen_addr = '[::]:50051'
    server.add_insecure_port(listen_addr)
    logging.info("Python grpc.aio 服务已在 %s 启动", listen_addr)

    await server.start()

    # 优雅退出信号处理
    loop = asyncio.get_running_loop()
    stop_event = asyncio.Event()

    for sig in (signal.SIGINT, signal.SIGTERM):
        loop.add_signal_handler(sig, stop_event.set)

    await stop_event.wait()
    logging.info("正在平滑关闭 gRPC 服务...")
    await server.stop(grace=5.0)
    logging.info("gRPC 服务已安全退出")

if __name__ == '__main__':
    logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
    asyncio.run(serve())

3. 现代异步客户端实战 (grpc.aio)

import asyncio
import logging
import grpc

import order_pb2
import order_pb2_grpc

async def main():
    # 建立异步长连接通道
    async with grpc.aio.insecure_channel('localhost:50051') as channel:
        stub = order_pb2_grpc.OrderServiceStub(channel)

        # 1. 发起一元调用 (携带超时控制)
        try:
            req = order_pb2.CreateOrderRequest(
                user_id="py-user-007",
                amount=88.5,
                item_ids=["book-ai", "keyboard"]
            )
            resp = await stub.CreateOrder(req, timeout=3.0)
            logging.info("创建订单返回: ID=%s, Status=%s", resp.order_id, resp.status)

            # 2. 消费服务端流
            track_req = order_pb2.TrackOrderRequest(order_id=resp.order_id)
            stream_call = stub.TrackOrder(track_req, timeout=10.0)

            async for event in stream_call:
                logging.info("[流更新] 状态: %s | 详情: %s", event.status, event.message)

        except grpc.aio.AioRpcError as rpc_err:
            logging.error("RPC 失败! 状态码: %s, 原因: %s", rpc_err.code(), rpc_err.details())

if __name__ == '__main__':
    logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
    asyncio.run(main())

4. 同步 vs 异步性能与技术选型

特性 同步多线程 (grpc.server) 异步协程 (grpc.aio.server)
底层模型 OS 操作系统线程池(ThreadPoolExecutor Python 事件循环(asyncio
内存开销 较高(每个线程占用 2MB~8MB 栈内存) 极低(每个协程仅需几百字节)
流式连接承载力 受线程数硬限制(通常 $\le 1000$) 轻松承载数万活跃长连接与流
适用场景 包含强 CPU 阻塞型计算(如 C 扩展、老旧同步 DB 驱动) 高频微服务通信、LLM 流式 Token 推送、WebSockets