Python 基础体系 · 第 89/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python gRPC:Protobuf、流式调用、Deadline、拦截器和错误模型

gRPC 是一种基于 HTTP/2 的远程过程调用(Remote Procedure Call,RPC)框架。它允许调用方使用本地函数般的接口调用远端服务,同时通过服务契约生成客户端和服务端代码。Python gRPC 通常使用 Protocol Buffers(Protobuf)描述消息和服务,再由代码生成器生成 Python 类型、序列化逻辑、客户端 Stub 和服务端基类。(grpc.io)

这套模型与 FastAPI 的 HTTP/JSON API 有明显不同:

  • FastAPI 的主要契约是 URL、HTTP 方法、JSON Schema 和 HTTP 状态码。
  • gRPC 的主要契约是 .proto 文件中的消息类型、服务名和 RPC 方法。
  • FastAPI 运行在 ASGI 应用生命周期之上。
  • gRPC Python 使用 grpc.Servergrpc.aio.Server 管理 HTTP/2 RPC 生命周期。

因此,gRPC 的核心问题不是“如何把一个函数暴露出去”,而是如何设计一套稳定的二进制契约,并正确处理流、截止时间、取消、跨服务错误和横切逻辑。


一、先建立 gRPC 的运行模型

一次最简单的 gRPC 调用可以抽象为:

客户端 StubChannelHTTP/2 RPC服务端 Handler\text{客户端 Stub} \rightarrow \text{Channel} \rightarrow \text{HTTP/2 RPC} \rightarrow \text{服务端 Handler}

其中:

  • Stub 是根据 .proto 生成的客户端代理;
  • Channel 表示客户端到服务端的连接抽象,负责连接、名称解析、负载均衡和流控等网络细节;
  • RPC 方法 是服务契约中声明的方法;
  • Handler 是服务端实现的方法;
  • Protobuf Message 是请求和响应在应用层使用的结构化对象;
  • Metadata 是不属于业务消息本身的附加键值对,例如认证信息、追踪 ID 和租户信息。

gRPC RPC 的结果不是单纯的“返回值”或“抛出 Python 异常”,而是一个状态:

RPC Result=(status code,status message,response,trailers)\text{RPC Result} = (\text{status code}, \text{status message}, \text{response}, \text{trailers})

成功时通常返回响应消息和 OK 状态;失败时返回 gRPC 状态码、错误描述以及可选的错误详情。这个错误模型独立于 Protobuf,即使底层数据格式不是 Protobuf,也可以使用 gRPC 的状态模型。(grpc.io)


二、Protobuf:消息契约和线格式

2.1 .proto 同时描述数据和服务

下面定义一个商品目录服务:

syntax = "proto3";

package catalog.v1;

message GetItemRequest {
  string id = 1;
}

message Item {
  string id = 1;
  string name = 2;
  int64 price_cent = 3;

  // 显式区分“没有备注”和“备注为空字符串”
  optional string note = 4;
}

message ListItemsRequest {
  repeated string ids = 1;
}

message UploadSummary {
  int32 accepted = 1;
}

service Catalog {
  rpc GetItem(GetItemRequest) returns (Item);

  rpc WatchItems(ListItemsRequest) returns (stream Item);

  rpc UploadItems(stream Item) returns (UploadSummary);

  rpc ChatItems(stream Item) returns (stream Item);
}

服务定义中的 stream 决定了 RPC 的通信方向:

请求 响应 类型
单个消息 单个消息 Unary-Unary
单个消息 消息流 Unary-Stream
消息流 单个消息 Stream-Unary
消息流 消息流 Stream-Stream

例如:

rpc GetItem(GetItemRequest) returns (Item);

是普通的一元 RPC;而:

rpc WatchItems(ListItemsRequest) returns (stream Item);

表示客户端发送一次请求,服务端可以连续返回多个 Item

stream 不是“把一个大数组拆成多次返回”的简单语法。它改变了 RPC 的生命周期:调用在第一个消息返回后仍然保持打开,直到服务端正常结束、客户端取消、Deadline 到期或发生传输错误。


2.2 字段编号是线协议的一部分

字段定义中的数字不是普通的序号,而是二进制线格式中的字段编号:

string id = 1;
string name = 2;
int64 price_cent = 3;

序列化后的消息不会依赖 Python 属性名来识别字段,而是依赖字段编号和字段类型。因此,下面这些变化通常是兼容的:

// 新增字段
optional string currency = 5;

旧客户端不会识别字段 5,但通常会跳过它;新客户端读取旧消息时,会使用字段默认值。

下面这些变化则可能破坏兼容性:

// 错误:复用旧字段编号
string currency = 3;

因为旧版本把字段 3 解释为 int64 price_cent,新版本却把它解释成字符串。删除字段时应保留编号:

message Item {
  reserved 3;
  reserved "price_cent";
}

这样可以让编译器阻止未来重新使用已经废弃的字段编号或名称。


2.3 Protobuf 不是 Python dataclass

生成的 Protobuf 消息具有自己的语义,不应把它当作普通字典或 dataclass

item = catalog_pb2.Item(
    id="sku-1",
    name="keyboard",
    price_cent=12900,
)

标量字段可以直接读取和赋值:

item.price_cent = 14900
print(item.price_cent)

但嵌套消息不能直接替换成普通 Python 对象:

message Warehouse {
  string name = 1;
}

message Item {
  string id = 1;
  Warehouse warehouse = 2;
}

下面的写法是不正确的:

item.warehouse = catalog_pb2.Warehouse(name="Hangzhou")

应当复制或修改嵌套消息:

item.warehouse.CopyFrom(
    catalog_pb2.Warehouse(name="Hangzhou")
)

或者直接修改其字段:

item.warehouse.name = "Hangzhou"

Python 生成代码中的消息类型由 Protobuf 运行时结合描述符生成,字段类型错误时会抛出 TypeError。(protobuf.dev)


2.4 默认值和字段存在性不是一回事

这是 Protobuf 最容易造成业务错误的地方之一。

proto3 中,普通标量字段默认不追踪存在性:

message UpdatePriceRequest {
  int64 price_cent = 1;
}

客户端无法区分以下两种请求:

1. 调用方没有传 price_cent
2. 调用方显式传入 price_cent = 0

两者读取出来都可能是:

request.price_cent == 0

如果业务上必须区分“未修改”和“修改为零”,应使用 optional

message UpdatePriceRequest {
  optional int64 price_cent = 1;
}

Python 中可以使用 HasField

request = catalog_pb2.UpdatePriceRequest()

print(request.HasField("price_cent"))  # False

request.price_cent = 0

print(request.HasField("price_cent"))  # True

清除字段:

request.ClearField("price_cent")
print(request.HasField("price_cent"))  # False

repeatedmap 字段没有这种存在性语义:空集合只能表示没有元素,不能表示“明确传入了一个空集合”和“没有传入该字段”的区别。消息类型字段和 oneof 则具有显式存在性。(protobuf.dev)

因此,下面的更新逻辑是有问题的:

if request.price_cent:
    update_price(request.price_cent)

因为价格为 0 时会被误判为“未提供”。

使用 optional 后应写成:

if request.HasField("price_cent"):
    update_price(request.price_cent)

三、生成 Python 代码并运行一个完整服务

3.1 安装依赖

在 Python 3.14 环境中,至少需要:

python -m venv .venv
source .venv/bin/activate

python -m pip install grpcio grpcio-tools

Windows PowerShell:

python -m venv .venv
.venv\Scripts\Activate.ps1

python -m pip install grpcio grpcio-tools

grpcio 提供运行时,grpcio-tools 提供 protoc 的 Python 调用入口和 gRPC Python 代码生成插件。实际安装时应确认所选 grpcio 版本已经提供当前平台和 CPython 3.14 的可用发行包;如果没有对应 wheel,安装可能退化为本地编译,从而引入编译器和 gRPC C 核心依赖问题。官方 Python 快速入门将 grpciogrpcio-tools 分开安装。(grpc.io)

目录结构:

project/
├── catalog.proto
├── server.py
├── client.py
└── generated/

生成代码:

python -m grpc_tools.protoc \
  -I. \
  --python_out=generated \
  --grpc_python_out=generated \
  catalog.proto

生成结果通常包括:

generated/
├── catalog_pb2.py
└── catalog_pb2_grpc.py

如果生成目录没有 __init__.py,可以添加一个空文件:

touch generated/__init__.py

在 Windows 上可以手动创建同名空文件。


3.2 异步服务端

下面的服务端实现包含一元调用、服务端流、客户端流和双向流:

# server.py
from __future__ import annotations

import asyncio
from collections.abc import AsyncIterator

import grpc

from generated import catalog_pb2
from generated import catalog_pb2_grpc


ITEMS = {
    "sku-1": catalog_pb2.Item(
        id="sku-1",
        name="Keyboard",
        price_cent=12900,
        note="机械键盘",
    ),
    "sku-2": catalog_pb2.Item(
        id="sku-2",
        name="Mouse",
        price_cent=5900,
    ),
}


class CatalogService(catalog_pb2_grpc.CatalogServicer):
    async def GetItem(
        self,
        request: catalog_pb2.GetItemRequest,
        context: grpc.aio.ServicerContext,
    ) -> catalog_pb2.Item:
        item = ITEMS.get(request.id)

        if item is None:
            await context.abort(
                grpc.StatusCode.NOT_FOUND,
                f"item {request.id!r} was not found",
            )

        # CopyFrom 或 CopyTo 等价语义都比直接返回可变共享对象更安全。
        return catalog_pb2.Item().CopyFrom(item)

    async def WatchItems(
        self,
        request: catalog_pb2.ListItemsRequest,
        context: grpc.aio.ServicerContext,
    ) -> AsyncIterator[catalog_pb2.Item]:
        for item_id in request.ids:
            if not context.is_active():
                return

            item = ITEMS.get(item_id)
            if item is None:
                continue

            await asyncio.sleep(0.1)
            yield item

    async def UploadItems(
        self,
        request_iterator: AsyncIterator[catalog_pb2.Item],
        context: grpc.aio.ServicerContext,
    ) -> catalog_pb2.UploadSummary:
        accepted = 0

        async for item in request_iterator:
            if not item.id or not item.name:
                await context.abort(
                    grpc.StatusCode.INVALID_ARGUMENT,
                    "id and name are required",
                )

            ITEMS[item.id] = catalog_pb2.Item().CopyFrom(item)
            accepted += 1

        return catalog_pb2.UploadSummary(accepted=accepted)

    async def ChatItems(
        self,
        request_iterator: AsyncIterator[catalog_pb2.Item],
        context: grpc.aio.ServicerContext,
    ) -> AsyncIterator[catalog_pb2.Item]:
        async for item in request_iterator:
            yield catalog_pb2.Item(
                id=item.id,
                name=f"received:{item.name}",
                price_cent=item.price_cent,
            )


async def serve() -> None:
    server = grpc.aio.server(
        maximum_concurrent_rpcs=100,
    )

    catalog_pb2_grpc.add_CatalogServicer_to_server(
        CatalogService(),
        server,
    )

    server.add_insecure_port("[::]:50051")
    await server.start()

    try:
        await server.wait_for_termination()
    finally:
        await server.stop(grace=5)


if __name__ == "__main__":
    asyncio.run(serve())

这里有几个关键点。

第一,grpc.aio.server() 创建的是 AsyncIO 服务端。服务方法使用 async def,服务端流方法使用 yield,客户端流方法接收异步迭代器。grpc.aio API 当前文档仍标记为实验性 API,因此应在项目中固定依赖版本并通过集成测试验证升级行为。(grpc.github.io)

第二,context.abort() 不只是设置错误状态,它会终止当前 RPC,并通过抛出内部异常把控制权交还给 gRPC 运行时。因此,abort() 后的普通返回语句不会执行。

第三,maximum_concurrent_rpcs=100 是服务端并发上限。超过限制的请求会得到 RESOURCE_EXHAUSTED,它与应用层限流不同:前者保护 gRPC 服务实例的并发能力,后者通常按用户、租户、IP 或业务资源执行更细粒度的配额控制。(grpc.github.io)


3.3 异步客户端

# client.py
from __future__ import annotations

import asyncio

import grpc

from generated import catalog_pb2
from generated import catalog_pb2_grpc


async def upload_items():
    async def requests():
        yield catalog_pb2.Item(
            id="sku-3",
            name="USB Cable",
            price_cent=1900,
        )
        yield catalog_pb2.Item(
            id="sku-4",
            name="USB Hub",
            price_cent=8900,
        )

    async with grpc.aio.insecure_channel("localhost:50051") as channel:
        stub = catalog_pb2_grpc.CatalogStub(channel)

        result = await stub.UploadItems(
            requests(),
            timeout=3.0,
        )

        print("accepted =", result.accepted)


async def main() -> None:
    async with grpc.aio.insecure_channel("localhost:50051") as channel:
        stub = catalog_pb2_grpc.CatalogStub(channel)

        item = await stub.GetItem(
            catalog_pb2.GetItemRequest(id="sku-1"),
            timeout=2.0,
        )
        print(item.name, item.price_cent)

        print("watch:")
        stream = stub.WatchItems(
            catalog_pb2.ListItemsRequest(
                ids=["sku-1", "sku-2", "missing"],
            ),
            timeout=2.0,
        )

        async for item in stream:
            print(" -", item.id, item.name)

        await upload_items()


if __name__ == "__main__":
    asyncio.run(main())

启动:

python server.py

另一个终端运行:

python client.py

预期输出类似:

Keyboard 12900
watch:
 - sku-1 Keyboard
 - sku-2 Mouse
accepted = 2

客户端应复用 Channel 和 Stub,而不是每次请求都创建一个新连接。Channel 包含连接管理和底层网络状态;频繁创建连接会增加握手、名称解析和资源管理开销。gRPC 官方性能建议也强调复用 Channel 和 Stub。(grpc.io)


四、四种流式调用到底改变了什么

4.1 Unary-Unary:请求和响应都是单个消息

response = await stub.GetItem(request, timeout=1.0)

它的生命周期是:

发送请求 → 等待响应 → RPC 结束

适合查询单个资源、执行一个明确的命令或返回有限结果。


4.2 Unary-Stream:服务端流

stream = stub.WatchItems(request, timeout=5.0)

async for item in stream:
    process(item)

生命周期是:

发送一个请求
    ↓
服务端返回零个或多个消息
    ↓
服务端发送最终状态

注意,async for 结束不一定意味着业务结果为空。它可能表示:

  1. 服务端正常发送完所有消息;
  2. 服务端返回了错误;
  3. 客户端取消了 RPC;
  4. Deadline 到期。

因此,完整错误处理应包围整个迭代过程:

try:
    async for item in stub.WatchItems(request, timeout=5.0):
        process(item)
except grpc.aio.AioRpcError as exc:
    print(exc.code(), exc.details())

不能只捕获创建 Stream 时的异常,因为流式 RPC 的错误经常发生在后续读取阶段。


4.3 Stream-Unary:客户端流

result = await stub.UploadItems(requests())

客户端持续发送多个请求消息,服务端最后只返回一个汇总结果。

这种模式适合批量上传、批量校验和批量写入。服务端不能因为收到了第一个消息就认为调用结束,必须等客户端发送完毕:

async for item in request_iterator:
    validate(item)

return summary

客户端迭代器抛出异常、连接断开或 Deadline 到期时,服务端的 async for 可能提前结束或抛出 RPC 相关异常。服务端写入数据库时应考虑“已经处理了一部分消息,但最终 RPC 失败”的情况。

例如:

客户端发送 item-1
服务端写入 item-1
客户端发送 item-2
服务端写入 item-2
Deadline 到期

此时客户端可能只看到 DEADLINE_EXCEEDED,但 item-1 和 item-2 可能已经生效。流式 RPC 不自动提供事务回滚。


4.4 Stream-Stream:双向流

call = stub.ChatItems(request_iterator)

async for response in call:
    process(response)

双向流允许客户端和服务端独立发送消息。常见实现是两个并发任务:

async def send(call):
    for item in input_items:
        await call.write(item)
    await call.done_writing()


async def receive(call):
    async for item in call:
        handle(item)

实际使用时也可以直接把异步迭代器传给 Stub:

responses = stub.ChatItems(requests())

async for response in responses:
    print(response)

双向流的核心不是“同时调用两个函数”,而是两条独立的消息方向:

sequenceDiagram
    participant C as 客户端
    participant G as gRPC/HTTP2
    participant S as 服务端

    C->>G: Item 1
    G->>S: Item 1
    S-->>G: Response 1
    G-->>C: Response 1

    C->>G: Item 2
    S-->>G: Response 2
    G-->>C: Response 2

    C->>G: 半关闭请求方向
    S-->>C: 最后一个响应
    S-->>C: OK trailers

客户端发送完请求并不一定关闭整个 RPC,只是关闭客户端到服务端的发送方向。服务端仍然可以继续发送响应,直到服务端结束或 RPC 被取消。


五、流式调用、背压和内存边界

背压(backpressure) 是指接收方处理速度较慢时,发送方不能无限制地继续缓存数据。

gRPC 的流控由 HTTP/2 和 gRPC 框架共同参与。应用层写入消息并不表示消息已经到达网络对端;消息可能仍在 gRPC 框架、操作系统或 HTTP/2 缓冲区中。接收方读取消息后,底层协议才会逐步释放发送方继续发送的能力。(grpc.io)

因此,下面的代码虽然不会一次性构造大列表,但仍可能产生过快的生产速度:

async def requests():
    for row in database_rows():
        yield to_proto(row)

如果 database_rows() 很快,而服务端处理很慢,流控最终会让发送方等待,但中间层仍可能持有一定数量的缓冲数据。

更明确的应用层控制方式是使用有界队列:

import asyncio


async def producer(queue: asyncio.Queue, rows):
    for row in rows:
        await queue.put(row)

    await queue.put(None)


async def request_iterator(queue: asyncio.Queue):
    while True:
        item = await queue.get()
        if item is None:
            return
        yield item

队列容量 N 给出了一个可解释的内存边界:

待发送消息数N\text{待发送消息数} \leq N

如果不设置上界,生产者可能因为数据库、文件或消息队列读取速度更快而在应用层积累大量对象。

流式调用还存在两个常见反例:

反例一:用服务端流返回数百万条记录

rpc ExportAll(ExportRequest) returns (stream Row);

这会让一个 RPC 长时间占用连接和服务端资源,而且流中途失败后,客户端通常只能知道“传输没有完成”,需要自行设计断点续传、游标或分页协议。

对于可重试、可恢复的批量导出,分页往往更容易运维:

message PageRequest {
  int32 page_size = 1;
  string page_token = 2;
}

message PageResponse {
  repeated Item items = 1;
  string next_page_token = 2;
}

反例二:双方都先写后读

如果客户端和服务端都试图大量写入,并且双方都在等待对方读取,可能形成死锁。官方流控说明明确指出,双向流在同步读写或手动流控场景下存在这种风险。(grpc.io)

双向流应当设计为:

  • 发送和接收并发进行;
  • 消息有界;
  • 明确定义半关闭行为;
  • 连接取消时停止生产;
  • 服务端检测 context.is_active()

六、Deadline:远程调用必须有时间边界

6.1 Timeout 和 Deadline 的区别

在 Python 调用中通常写:

await stub.GetItem(request, timeout=2.0)

这里的 timeout=2.0 是客户端 API 使用的相对时长,表示本次 RPC 最多等待约两秒。

从分布式系统语义看,Deadline 更接近一个绝对截止时刻:

D=t0+ΔD = t_0 + \Delta

其中:

  • t0t_0 是调用开始时间;
  • Δ\Delta 是调用允许的总时长;
  • DD 是必须完成 RPC 的截止时间。

如果服务 A 调用服务 B,再由 B 调用服务 C,不能让每一跳都重新获得完整的两秒:

A → B:2 秒
B → C:2 秒

这样整个请求可能等待四秒甚至更久,超过用户请求的预算。

正确的预算应该是:

ΔCDtnow\Delta_C \leq D - t_{\text{now}}

例如:

用户总预算:2.0 秒
A 处理和网络耗时:0.3 秒
B 调用 C 时剩余预算:约 1.7 秒

服务端可以通过:

remaining = context.time_remaining()

读取当前 RPC 剩余时间。没有设置 Deadline 时,该值可能为 None;设置后,它是一个非负的剩余秒数。(grpc.github.io)


6.2 Deadline 到期时会发生什么

gRPC 默认不自动设置 Deadline,因此如果客户端不传 timeout,调用可能无限等待。Deadline 到期后,客户端通常看到:

grpc.StatusCode.DEADLINE_EXCEEDED

服务端会自动取消对应调用。(grpc.io)

客户端:

try:
    response = await stub.GetItem(
        catalog_pb2.GetItemRequest(id="sku-1"),
        timeout=0.01,
    )
except grpc.aio.AioRpcError as exc:
    if exc.code() == grpc.StatusCode.DEADLINE_EXCEEDED:
        print("request timed out")

服务端长任务应主动检查取消状态:

async def expensive_work(context):
    for step in range(100):
        if not context.is_active():
            return

        await asyncio.sleep(0.05)
        do_one_step(step)

仅仅让客户端超时,并不意味着服务端的业务代码会立即停止。服务端可能已经把数据写入数据库、发送消息或调用下游。Deadline 是取消信号和结果边界,不是自动事务回滚。


6.3 给下游调用传递剩余 Deadline

假设服务端收到 RPC 后还要调用另一个 gRPC 服务:

remaining = context.time_remaining()

if remaining is None:
    # 没有上游 Deadline 时,也不应无限等待下游。
    downstream_timeout = 3.0
else:
    # 留出少量本地收尾时间,避免下游恰好耗尽全部预算。
    downstream_timeout = max(0.001, remaining - 0.05)

response = await downstream_stub.GetItem(
    request,
    timeout=downstream_timeout,
)

这里的 0.05 不是规范要求,而是应用层预算策略。它需要根据序列化、日志、数据库提交和响应发送时间调整,不能当作普适性能数字。

Deadline 与重试必须一起设计。对于:

客户端请求 → 服务端执行写操作 → 客户端 Deadline 到期

客户端无法仅凭 DEADLINE_EXCEEDED 判断服务端是否已经提交成功。因此,非幂等写操作不能因为超时就盲目重试。应通过幂等键、请求 ID 或服务端去重表解决:

message CreateOrderRequest {
  string idempotency_key = 1;
  string user_id = 2;
  int64 amount_cent = 3;
}

服务端以 idempotency_key 唯一约束确保重复请求不会创建多个订单。


七、Metadata:请求之外的上下文通道

Metadata 是与 RPC 关联的键值对,底层通过 HTTP/2 headers 和 trailers 传递。它适合承载认证凭据、追踪信息、租户标识和少量诊断信息,而不适合承载业务主体。Metadata 的键不区分大小写,应用自定义键不应使用 grpc- 前缀;服务器还可能限制请求头大小。(grpc.io)

客户端发送 metadata:

response = await stub.GetItem(
    request,
    metadata=(
        ("authorization", "Bearer token"),
        ("x-request-id", "req-123"),
    ),
    timeout=2.0,
)

服务端读取:

metadata = dict(context.invocation_metadata())

request_id = metadata.get("x-request-id")
authorization = metadata.get("authorization")

如果服务端返回错误,可以附加 trailers:

await context.abort(
    grpc.StatusCode.RESOURCE_EXHAUSTED,
    "quota exceeded",
    trailing_metadata=(
        ("retry-after-ms", "1000"),
    ),
)

retry-after-ms 只是应用约定,不是 gRPC 标准状态字段。客户端需要把它当作不可信输入进行解析和上限保护。


八、拦截器:在 RPC 边界插入横切逻辑

拦截器(Interceptor) 是在 RPC 调用进入实际 Handler 前后执行的扩展点,适合实现:

  • 认证信息注入;
  • 请求 ID;
  • 日志;
  • 指标;
  • 统一异常映射;
  • 访问控制;
  • Trace 上下文传播。

拦截器不应该替代业务 Handler,也不应该隐藏会改变业务语义的重试。

8.1 AsyncIO 客户端拦截器

# interceptors.py
from __future__ import annotations

from uuid import uuid4

import grpc


class RequestIdInterceptor(
    grpc.aio.UnaryUnaryClientInterceptor,
):
    async def intercept_unary_unary(
        self,
        continuation,
        client_call_details,
        request,
    ):
        metadata = list(client_call_details.metadata or ())
        metadata.append(("x-request-id", uuid4().hex))

        new_details = client_call_details._replace(
            metadata=metadata,
        )

        return await continuation(new_details, request)

创建 Channel:

interceptor = RequestIdInterceptor()

async with grpc.aio.insecure_channel(
    "localhost:50051",
    interceptors=(interceptor,),
) as channel:
    stub = catalog_pb2_grpc.CatalogStub(channel)
    response = await stub.GetItem(
        catalog_pb2.GetItemRequest(id="sku-1"),
        timeout=2.0,
    )

AsyncIO 客户端拦截器按照 RPC 类型区分,包括:

  • UnaryUnaryClientInterceptor
  • UnaryStreamClientInterceptor
  • StreamUnaryClientInterceptor
  • StreamStreamClientInterceptor

拦截器必须调用 continuation 才会继续执行后续拦截器或真正的 RPC。当前 Python AsyncIO 拦截器 API 在官方文档中标为实验性接口,且流式 RPC 的拦截器返回值可能是 Call 对象,也可能是异步迭代器,因此不能把一元 RPC 拦截器代码机械复制到流式 RPC。(grpc.github.io)

8.2 服务器端拦截器

同步服务端:

server = grpc.server(
    thread_pool,
    interceptors=(server_interceptor,),
)

异步服务端:

server = grpc.aio.server(
    interceptors=(server_interceptor,),
)

服务器端拦截器的入口是:

async def intercept_service(
    self,
    continuation,
    handler_call_details,
):
    return await continuation(handler_call_details)

它接收的方法路径、请求 metadata 等信息,并可以:

  1. 拒绝请求;
  2. 调用后续拦截器;
  3. 返回包装后的 Handler;
  4. 使用 contextvars 向下游传递请求上下文。

服务端拦截器适合做统一认证和指标,但要注意两个边界:

  • 它不能直接访问业务请求对象,除非进一步包装具体 Handler;
  • 它不能保证拦截器和 Handler 在同一个线程中执行,不能依赖线程局部变量。官方 Python API 明确说明,服务端拦截器和 Handler 之间可使用 contextvars 传递状态。(grpc.github.io)

九、错误模型:不要把所有失败都变成 INTERNAL

9.1 标准状态码

gRPC 的状态码是跨语言协议的一部分。常用状态码包括:

状态码 典型含义
INVALID_ARGUMENT 请求参数本身非法
UNAUTHENTICATED 缺少或无法验证身份凭据
PERMISSION_DENIED 已识别调用方,但没有权限
NOT_FOUND 资源不存在
ALREADY_EXISTS 创建资源时资源已存在
FAILED_PRECONDITION 当前系统状态不满足操作条件
ABORTED 并发冲突或事务中止
RESOURCE_EXHAUSTED 配额、容量或并发资源耗尽
UNAVAILABLE 服务暂时不可用,可能适合重试
DEADLINE_EXCEEDED Deadline 到期
CANCELLED 调用方取消
INTERNAL 服务内部不变量被破坏
UNIMPLEMENTED 方法未实现或不支持

其中 INVALID_ARGUMENTNOT_FOUNDALREADY_EXISTSFAILED_PRECONDITIONABORTEDOUT_OF_RANGEDATA_LOSS 通常只由应用代码主动返回,gRPC 库不会凭空生成这些业务状态码。(grpc.io)


9.2 FAILED_PRECONDITIONABORTEDUNAVAILABLE 的区别

这三个状态码经常被混用。

假设库存扣减失败:

库存为负,必须先修复数据

更接近:

grpc.StatusCode.FAILED_PRECONDITION

如果是乐观锁版本冲突:

客户端使用 version=10 更新,但数据库当前已经是 version=11

更接近:

grpc.StatusCode.ABORTED

如果数据库连接池暂时不可用:

稍后重试同一个 RPC 可能成功

更接近:

grpc.StatusCode.UNAVAILABLE

官方状态码指南给出的判断原则是:

  • 如果只需重试当前调用,考虑 UNAVAILABLE
  • 如果需要在更高层重新执行一组操作,考虑 ABORTED
  • 如果必须先修复系统状态,考虑 FAILED_PRECONDITION。(grpc.io)

9.3 Python 客户端捕获错误

try:
    item = await stub.GetItem(
        catalog_pb2.GetItemRequest(id="missing"),
        timeout=2.0,
    )
except grpc.aio.AioRpcError as exc:
    print("code =", exc.code())
    print("details =", exc.details())
    print("trailers =", exc.trailing_metadata())

    if exc.code() == grpc.StatusCode.NOT_FOUND:
        print("资源不存在")
    elif exc.code() == grpc.StatusCode.UNAVAILABLE:
        print("服务暂时不可用")
    elif exc.code() == grpc.StatusCode.DEADLINE_EXCEEDED:
        print("超过截止时间")
    else:
        raise

details() 是面向人的字符串,不应该被客户端解析成结构化业务协议。下面这种做法不稳定:

if "quota" in exc.details():
    retry()

错误文本可能在服务端升级后变化,也可能被本地化或重新包装。重试判断应使用状态码和明确的错误详情类型。


十、丰富错误模型:把结构化详情放进 Status

标准错误模型包含状态码和字符串描述。如果客户端还需要得到字段级校验错误、资源名称、重试时间或帮助链接,可以使用丰富错误模型。

例如,服务端返回字段校验错误:

from google.protobuf import any_pb2
from google.rpc import error_details_pb2
from google.rpc import status_pb2
from grpc_status import rpc_status


async def abort_invalid_request(context, field: str):
    violation = error_details_pb2.BadRequest.FieldViolation(
        field=field,
        description="must not be empty",
    )

    detail = any_pb2.Any()
    detail.Pack(violation)

    status = status_pb2.Status(
        code=grpc.StatusCode.INVALID_ARGUMENT.value[0],
        message="request validation failed",
        details=[detail],
    )

    await context.abort_with_status(
        rpc_status.to_status(status)
    )

需要额外安装:

python -m pip install grpcio-status googleapis-common-protos

客户端解析:

from grpc_status import rpc_status

try:
    await stub.GetItem(request, timeout=2.0)
except grpc.aio.AioRpcError as exc:
    status = rpc_status.from_call(exc)

    if status is not None:
        print(status.code)
        print(status.message)

        for detail in status.details:
            if detail.Is(error_details_pb2.BadRequest.DESCRIPTOR):
                bad_request = error_details_pb2.BadRequest()
                detail.Unpack(bad_request)

                for violation in bad_request.field_violations:
                    print(
                        violation.field,
                        violation.description,
                    )

abort_with_status 当前在 Python gRPC API 中属于实验性能力,因此应固定 grpcio-statusgrpcio 的兼容版本。(grpc.github.io)

结构化错误的价值在于让客户端能稳定地做出行为:

INVALID_ARGUMENT + BadRequest
    → 修正输入,不重试

RESOURCE_EXHAUSTED + RetryInfo
    → 等待指定时间后再尝试

UNAVAILABLE
    → 根据幂等性和重试预算决定是否重试

NOT_FOUND
    → 展示资源不存在,不做无限重试

十一、错误边界:传输错误、业务错误和部分成功

不要把业务错误编码成正常响应:

message GetItemResponse {
  bool ok = 1;
  string error = 2;
  Item item = 3;
}

这种设计的问题是:即使业务失败,gRPC 层仍然看到 OK。监控、重试中间件和调用方无法通过标准 RPC 状态判断请求是否失败。

更自然的定义是:

rpc GetItem(GetItemRequest) returns (Item);

业务失败使用 gRPC 状态:

await context.abort(
    grpc.StatusCode.NOT_FOUND,
    "item was not found",
)

但这不意味着所有信息都应该放进 gRPC 状态。大批量校验结果属于业务结果,应放在响应消息中:

message BatchResult {
  int32 accepted = 1;
  int32 rejected = 2;
  repeated RejectedItem rejected_items = 3;
}

可以使用如下划分:

  • RPC 状态:调用是否完成、是否可重试、调用方是否有权限;
  • 响应消息:业务处理结果和部分成功明细;
  • Metadata:认证、追踪和少量跨层信息;
  • 日志和 Trace:内部堆栈、数据库错误和诊断上下文。

服务器端不应把数据库连接字符串、Python 堆栈或内部表名直接放入 details()。客户端需要可见的错误信息应经过分类和脱敏。


十二、gRPC 与 Python 异步应用的生命周期

12.1 gRPC 不是 ASGI 应用

ASGI 定义的是异步 Python 应用与服务器之间的接口规范;FastAPI 是构建在 ASGI 生态上的 HTTP Web 框架。gRPC Python 使用自己的 Server、Channel 和 RPC Handler 模型,不会因为服务方法是 async def 就自动成为 ASGI 应用。

如果一个进程同时提供 FastAPI 和 gRPC,常见部署方式是分别监听端口:

:8000  FastAPI / HTTP
:50051 gRPC

两者可以共享:

  • 数据库访问层;
  • 领域服务;
  • 配置;
  • 认证逻辑;
  • Trace 和指标组件。

但不应直接共享 HTTP Handler 和 gRPC Handler,因为两套协议的错误、取消、流和生命周期语义不同。

12.2 优雅关闭

服务端关闭时应停止接收新请求,并给正在执行的 RPC 留出完成时间:

try:
    await server.wait_for_termination()
finally:
    await server.stop(grace=5)

grace=5 表示等待正在进行的 RPC 最多约五秒;它不是保证所有业务都能完成的事务时间。如果服务方法阻塞在不可取消的同步 I/O 上,优雅关闭也可能无法及时完成。

Python AsyncIO 中还要注意取消传播:

async def handler(context):
    task = asyncio.create_task(expensive_work())

    try:
        return await task
    except asyncio.CancelledError:
        task.cancel()
        raise

如果收到客户端取消后仍让后台任务继续写数据库或发送消息,就可能造成“客户端已失败、服务端仍完成写入”的结果。是否允许这种行为,必须由接口的幂等性和业务事务设计决定。


十三、同步 gRPC、grpc.aio 和性能取舍

同步 API 适合传统线程池服务:

server = grpc.server(
    futures.ThreadPoolExecutor(max_workers=20)
)

异步 API 适合已经采用 AsyncIO 的服务:

server = grpc.aio.server()

两者不能只通过把函数声明改成 async def 互换。同步 Stub、同步 Server、同步拦截器和异步对象属于不同 API 栈。

Python gRPC 的流式 RPC 与其他语言相比有额外性能特征。官方性能文档指出,Python 的流式调用可能为收发消息创建额外线程,因此通常比一元 RPC 更慢;使用 AsyncIO 可能改善这一点。官方还建议复用 Channel 和 Stub,并提醒长流无法在已经建立后重新进行负载均衡,流失败也更难调试。(grpc.io)

因此,不应为了“看起来更实时”而把所有接口设计成双向流:

  • 一元 RPC 更容易超时、重试和观测;
  • 服务端流适合连续结果,但需要处理半途失败;
  • 客户端流适合批量上传,但需要定义部分成功;
  • 双向流适合会话型协议,但需要处理并发读写、取消和状态同步。

十四、生产排查路径

当一次 gRPC 调用失败时,可以按以下因果链排查。

14.1 先看客户端状态码

except grpc.aio.AioRpcError as exc:
    print(exc.code())
    print(exc.details())

状态码先决定问题类别:

UNAUTHENTICATED → 凭据或认证链路
PERMISSION_DENIED → 授权策略
INVALID_ARGUMENT → 请求内容
NOT_FOUND → 资源定位
DEADLINE_EXCEEDED → 时间预算、下游慢或连接问题
UNAVAILABLE → 服务发现、连接、实例或网络
RESOURCE_EXHAUSTED → 配额、并发或流控
INTERNAL → 服务端缺陷或协议处理异常

14.2 再看请求阶段

需要区分:

连接前失败
请求发送前失败
请求已发送但没有响应
已经收到部分流消息
服务端完成业务但响应未送达

尤其是 DEADLINE_EXCEEDEDUNAVAILABLE,不代表服务端一定没有执行。对于写操作,必须通过服务端日志、幂等键或业务查询确认结果,而不能仅依赖客户端异常。

14.3 最后关联 Request ID 和 Trace

客户端通过拦截器注入:

x-request-id: req-abc

服务端日志至少记录:

request_id
rpc_method
grpc_code
deadline_remaining
peer
duration_ms

对于流式 RPC,还应记录:

messages_received
messages_sent
stream_duration
cancelled
last_message_id

因为一条“流失败”信息通常不足以判断是第一个消息失败,还是已经发送了数千条消息后连接中断。


十五、设计检查:把协议语义写进接口

一个可维护的 gRPC 接口至少要明确以下问题:

  1. 标量字段是否需要区分“未提供”和“默认值”;
  2. 字段编号是否已经分配并禁止复用;
  3. RPC 是一元调用还是流式调用;
  4. 流中途失败时,客户端如何知道哪些消息已经生效;
  5. 调用是否幂等;
  6. DEADLINE_EXCEEDED 后是否允许重试;
  7. UNAVAILABLE 是否真的代表当前操作可安全重试;
  8. 错误使用标准状态码还是结构化错误详情;
  9. 认证和 Trace 信息放在哪些 metadata 中;
  10. 服务端取消时,数据库和消息队列操作如何停止或收敛;
  11. FastAPI、HTTP API 和 gRPC API 是否共享领域逻辑但隔离协议适配层;
  12. Channel、连接和服务器是否在进程生命周期内复用并优雅关闭。

gRPC 的优势来自契约、生成代码、HTTP/2 流和跨语言状态模型的组合,而不是来自“把 HTTP 换成二进制”。只有把 Protobuf 的兼容规则、流式调用的状态、Deadline 的预算传播、拦截器的边界和错误码的行为一起设计,RPC 接口才真正具备微服务之间长期演进所需要的确定性。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。