返回博客

API 网关:统一入口的路由、限流与协议转换

coomia-dip API 网关作为平台的统一入口,实现请求路由(gRPC/REST 双协议)、认证鉴权、速率限流、协议转换(REST-to-gRPC)、负载均衡和请求/响应变换。网关基于 Envoy + 自定义 Python 过滤器构建,支持声明式路由配置、动态限流策略和 gRPC-Web 转码。本文从网关架构、路由设计、限流策略、协议转换、安全加固到高可用部署,完整解析 API 网关方案。

Coomia发布于 2025年10月6日9 分钟阅读
分享本文Twitter / X

系列:S6 平台工程 · 第 22 篇 | 难度:高级 | 阅读时间:18 分钟

API 网关:统一入口的路由、限流与协议转换

#TL;DR

coomia-dip API 网关作为平台的统一入口,实现请求路由(gRPC/REST 双协议)、认证鉴权、速率限流、协议转换(REST-to-gRPC)、负载均衡和请求/响应变换。网关基于 Envoy + 自定义 Python 过滤器构建,支持声明式路由配置、动态限流策略和 gRPC-Web 转码。本文从网关架构、路由设计、限流策略、协议转换、安全加固到高可用部署,完整解析 API 网关方案。

#1. 网关架构

#1.1 系统架构

Code
┌─────────── 外部流量 ─────────────┐
│  REST Clients │ gRPC Clients │ Web │
└───────────────┬──────────────────┘
                │
┌───────────────▼──────────────────┐
│          API Gateway              │
│  ┌──────────────────────────┐    │
│  │     Load Balancer         │    │
│  └────────────┬─────────────┘    │
│  ┌────────────▼─────────────┐    │
│  │    TLS Termination        │    │
│  └────────────┬─────────────┘    │
│  ┌────────────▼─────────────┐    │
│  │    Authentication         │    │
│  │    (JWT/OAuth2/API Key)   │    │
│  └────────────┬─────────────┘    │
│  ┌────────────▼─────────────┐    │
│  │    Rate Limiting          │    │
│  │    (Token Bucket)         │    │
│  └────────────┬─────────────┘    │
│  ┌────────────▼─────────────┐    │
│  │    Protocol Translation   │    │
│  │    (REST ↔ gRPC)          │    │
│  └────────────┬─────────────┘    │
│  ┌────────────▼─────────────┐    │
│  │    Request Router         │    │
│  │    (Service Discovery)    │    │
│  └────────────┬─────────────┘    │
└───────────────┼──────────────────┘
      ┌─────────┼─────────┐
      │         │         │
  ┌───▼──┐  ┌──▼───┐  ┌──▼───────┐
  │Ctrl  │  │Data  │  │Intel     │
  │Layer │  │Layer │  │Layer     │
  └──────┘  └──────┘  └──────────┘

#1.2 对标 Palantir Foundry

能力Palantir Foundrycoomia-dip
协议支持RESTREST + gRPC 双协议
协议转换无需REST-to-gRPC 转码
认证OAuth2JWT + OAuth2 + API Key
限流内置多维度令牌桶
gRPC-Web不支持完整支持

#2. 路由设计

#2.1 声明式路由配置

Python
class RouteConfig(BaseModel):
    """路由配置模型"""

    routes: list[Route]

class Route(BaseModel):
    """路由规则"""
    name: str
    match: RouteMatch
    upstream: UpstreamConfig
    policies: list[str] = Field(default_factory=list)
    transform: TransformConfig | None = None

class RouteMatch(BaseModel):
    """路由匹配条件"""
    prefix: str | None = None
    path: str | None = None
    method: str | None = None
    headers: dict[str, str] = Field(default_factory=dict)

class UpstreamConfig(BaseModel):
    """上游服务配置"""
    service: str
    port: int
    protocol: str = "grpc"  # grpc | http
    timeout_ms: int = 30000
    retry_policy: RetryPolicy | None = None


# 路由配置示例
GATEWAY_ROUTES = RouteConfig(routes=[
    # REST API → Ontology Service (gRPC)
    Route(
        name="ontology-objects",
        match=RouteMatch(prefix="/api/v1/objects"),
        upstream=UpstreamConfig(
            service="ontology-service", port=9090, protocol="grpc",
        ),
        policies=["auth", "rate-limit", "audit"],
        transform=TransformConfig(
            type="rest-to-grpc",
            service="OntologyService",
        ),
    ),

    # REST API → Schema Registry (gRPC)
    Route(
        name="schema-registry",
        match=RouteMatch(prefix="/api/v1/schemas"),
        upstream=UpstreamConfig(
            service="schema-registry", port=9091, protocol="grpc",
        ),
        policies=["auth", "rate-limit"],
    ),

    # gRPC 直通
    Route(
        name="grpc-passthrough",
        match=RouteMatch(prefix="/onto."),
        upstream=UpstreamConfig(
            service="ontology-service", port=9090, protocol="grpc",
        ),
        policies=["auth"],
    ),

    # gRPC-Web 转码
    Route(
        name="grpc-web",
        match=RouteMatch(
            prefix="/onto.",
            headers={"content-type": "application/grpc-web"},
        ),
        upstream=UpstreamConfig(
            service="ontology-service", port=9090, protocol="grpc",
        ),
        policies=["auth", "cors"],
        transform=TransformConfig(type="grpc-web"),
    ),

    # 健康检查(无认证)
    Route(
        name="health",
        match=RouteMatch(path="/health"),
        upstream=UpstreamConfig(
            service="gateway", port=8080, protocol="http",
        ),
        policies=[],
    ),
])

#2.2 动态路由发现

Python
class ServiceDiscovery:
    """服务发现 - 从 Kubernetes 或 Consul 获取上游服务地址"""

    async def resolve(self, service_name: str) -> list[Endpoint]:
        """解析服务地址"""
        if self._mode == "kubernetes":
            endpoints = await self._k8s_client.list_endpoints(
                namespace="onto-system",
                label_selector=f"app={service_name}",
            )
            return [
                Endpoint(address=ep.ip, port=ep.port, weight=1)
                for ep in endpoints
            ]
        elif self._mode == "static":
            return self._static_endpoints.get(service_name, [])

    async def watch(self, service_name: str, callback):
        """监听服务变更"""
        async for event in self._k8s_client.watch_endpoints(
            namespace="onto-system",
            label_selector=f"app={service_name}",
        ):
            endpoints = self._parse_endpoints(event)
            await callback(service_name, endpoints)

#3. 速率限流

#3.1 多维度令牌桶

Python
class RateLimiter:
    """多维度速率限流器"""

    def __init__(self):
        self._limiters: dict[str, TokenBucket] = {}
        self._redis: aioredis.Redis | None = None

    async def check_rate_limit(
        self,
        request: GatewayRequest,
    ) -> RateLimitResult:
        """检查速率限制"""
        dimensions = self._extract_dimensions(request)

        for dim_key, limit_config in self._get_applicable_limits(dimensions):
            bucket = await self._get_or_create_bucket(dim_key, limit_config)
            allowed = await bucket.consume(1)

            if not allowed:
                return RateLimitResult(
                    allowed=False,
                    limit=limit_config.requests_per_second,
                    remaining=0,
                    retry_after=bucket.next_available_time(),
                    dimension=dim_key,
                )

        return RateLimitResult(allowed=True)

    def _extract_dimensions(self, request: GatewayRequest) -> dict:
        return {
            "global": "global",
            "per_user": request.user_id,
            "per_ip": request.client_ip,
            "per_api": request.path,
            "per_object_type": request.headers.get("x-object-type", ""),
        }


class TokenBucket:
    """令牌桶实现(Redis 分布式)"""

    def __init__(
        self,
        redis: aioredis.Redis,
        key: str,
        rate: float,
        capacity: int,
    ):
        self._redis = redis
        self._key = key
        self._rate = rate        # 令牌填充速率(个/秒)
        self._capacity = capacity  # 桶容量

    async def consume(self, tokens: int = 1) -> bool:
        """尝试消费令牌"""
        lua_script = """
        local key = KEYS[1]
        local rate = tonumber(ARGV[1])
        local capacity = tonumber(ARGV[2])
        local now = tonumber(ARGV[3])
        local requested = tonumber(ARGV[4])

        local data = redis.call('HMGET', key, 'tokens', 'last_time')
        local tokens = tonumber(data[1]) or capacity
        local last_time = tonumber(data[2]) or now

        local elapsed = now - last_time
        tokens = math.min(capacity, tokens + elapsed * rate)

        if tokens >= requested then
            tokens = tokens - requested
            redis.call('HMSET', key, 'tokens', tokens, 'last_time', now)
            redis.call('EXPIRE', key, math.ceil(capacity / rate) * 2)
            return 1
        else
            redis.call('HMSET', key, 'tokens', tokens, 'last_time', now)
            return 0
        end
        """
        result = await self._redis.eval(
            lua_script, 1, self._key,
            self._rate, self._capacity, time.time(), tokens,
        )
        return bool(result)

#3.2 限流策略配置

Python
RATE_LIMIT_POLICIES = {
    "global": RateLimitConfig(
        requests_per_second=10000,
        burst=15000,
    ),
    "per_user": RateLimitConfig(
        requests_per_second=100,
        burst=200,
    ),
    "per_ip": RateLimitConfig(
        requests_per_second=50,
        burst=100,
    ),
    "per_api": {
        "/api/v1/objects/*/query": RateLimitConfig(
            requests_per_second=20,
            burst=50,
        ),
        "/api/v1/actions/*/execute": RateLimitConfig(
            requests_per_second=10,
            burst=20,
        ),
    },
}

#4. 协议转换

#4.1 REST-to-gRPC 转码

Python
class RestToGrpcTranscoder:
    """REST-to-gRPC 协议转换器"""

    def __init__(self, proto_descriptors: dict):
        self._descriptors = proto_descriptors

    async def transcode(
        self,
        http_request: HTTPRequest,
        route: Route,
    ) -> GrpcRequest:
        """将 REST 请求转换为 gRPC 请求"""
        # 1. 路径匹配到 gRPC 方法
        grpc_method = self._resolve_method(
            http_request.method,
            http_request.path,
            route.transform.service,
        )

        # 2. 构建 Protobuf 请求
        request_message = self._build_request_message(
            grpc_method,
            path_params=http_request.path_params,
            query_params=http_request.query_params,
            body=http_request.body,
        )

        return GrpcRequest(
            method=grpc_method,
            message=request_message,
            metadata=self._convert_headers(http_request.headers),
        )

    async def reverse_transcode(
        self,
        grpc_response: GrpcResponse,
    ) -> HTTPResponse:
        """将 gRPC 响应转换为 REST 响应"""
        return HTTPResponse(
            status_code=self._grpc_to_http_status(grpc_response.code),
            body=MessageToDict(grpc_response.message),
            headers=self._convert_metadata(grpc_response.metadata),
        )

    # REST 路径到 gRPC 方法的映射
    PATH_METHOD_MAP = {
        ("GET", "/api/v1/objects/{type}/{id}"): "OntologyService/GetObject",
        ("GET", "/api/v1/objects/{type}"): "OntologyService/ListObjects",
        ("POST", "/api/v1/objects/{type}"): "OntologyService/CreateObject",
        ("PUT", "/api/v1/objects/{type}/{id}"): "OntologyService/UpdateObject",
        ("DELETE", "/api/v1/objects/{type}/{id}"): "OntologyService/DeleteObject",
        ("POST", "/api/v1/objects/{type}/query"): "OntologyService/QueryObjects",
        ("POST", "/api/v1/actions/{type}/execute"): "ActionService/ExecuteAction",
    }

#4.2 gRPC-Web 支持

Python
class GrpcWebFilter:
    """gRPC-Web 协议过滤器"""

    async def handle(self, request: HTTPRequest) -> HTTPResponse:
        content_type = request.headers.get("content-type", "")

        if "application/grpc-web-text" in content_type:
            # Base64 编码的 gRPC-Web
            body = base64.b64decode(request.body)
        elif "application/grpc-web" in content_type:
            # 二进制 gRPC-Web
            body = request.body
        else:
            return HTTPResponse(status_code=415, body={"error": "Unsupported media type"})

        # 转换为标准 gRPC 请求
        grpc_request = self._grpc_web_to_grpc(body, request.headers)

        # 调用上游 gRPC 服务
        grpc_response = await self._forward_to_upstream(grpc_request)

        # 转换回 gRPC-Web 响应
        return self._grpc_to_grpc_web(grpc_response, content_type)

#5. 安全加固

#5.1 认证过滤器

Python
class AuthenticationFilter:
    """认证过滤器"""

    STRATEGIES = {
        "bearer": BearerTokenStrategy(),
        "api_key": APIKeyStrategy(),
        "oauth2": OAuth2Strategy(),
    }

    async def authenticate(self, request: GatewayRequest) -> AuthResult:
        auth_header = request.headers.get("authorization", "")

        if auth_header.startswith("Bearer "):
            return await self.STRATEGIES["bearer"].authenticate(auth_header)
        elif "x-api-key" in request.headers:
            return await self.STRATEGIES["api_key"].authenticate(
                request.headers["x-api-key"],
            )
        else:
            return AuthResult(authenticated=False, reason="No credentials provided")


class BearerTokenStrategy:
    async def authenticate(self, auth_header: str) -> AuthResult:
        token = auth_header.replace("Bearer ", "")

        try:
            payload = jwt.decode(
                token,
                self._public_key,
                algorithms=["RS256"],
                audience="coomia-dip",
            )
            return AuthResult(
                authenticated=True,
                user_id=payload["sub"],
                roles=payload.get("roles", []),
                claims=payload,
            )
        except jwt.ExpiredSignatureError:
            return AuthResult(authenticated=False, reason="Token expired")
        except jwt.InvalidTokenError as e:
            return AuthResult(authenticated=False, reason=str(e))

#5.2 CORS 配置

Python
CORS_CONFIG = CORSConfig(
    allowed_origins=["https://console.example.com"],
    allowed_methods=["GET", "POST", "PUT", "DELETE", "OPTIONS"],
    allowed_headers=["Authorization", "Content-Type", "X-Request-ID"],
    exposed_headers=["X-Request-ID", "X-RateLimit-Remaining"],
    max_age=3600,
    allow_credentials=True,
)

#6. 测试策略

Python
class TestAPIGateway:
    async def test_rest_to_grpc_transcoding(self):
        response = await http_client.get("/api/v1/objects/Employee/emp-001")
        assert response.status_code == 200
        assert response.json()["id"] == "emp-001"

    async def test_rate_limiting(self):
        for _ in range(100):
            await http_client.get("/api/v1/objects/Employee/emp-001")
        response = await http_client.get("/api/v1/objects/Employee/emp-001")
        assert response.status_code == 429  # Too Many Requests

    async def test_authentication_required(self):
        response = await http_client.get(
            "/api/v1/objects/Employee/emp-001",
            headers={},  # 无认证头
        )
        assert response.status_code == 401

    async def test_grpc_passthrough(self):
        channel = grpc.aio.insecure_channel("gateway:8443")
        stub = OntologyServiceStub(channel)
        response = await stub.GetObject(GetObjectRequest(
            object_type="Employee", object_id="emp-001",
        ))
        assert response.object.id == "emp-001"

    async def test_health_check_no_auth(self):
        response = await http_client.get("/health")
        assert response.status_code == 200

#7. 高可用部署

#7.1 部署架构

YAML
apiVersion: apps/v1
kind: Deployment
metadata:
  name: api-gateway
spec:
  replicas: 3
  strategy:
    type: RollingUpdate
    rollingUpdate:
      maxUnavailable: 1
      maxSurge: 1
  template:
    spec:
      affinity:
        podAntiAffinity:
          requiredDuringSchedulingIgnoredDuringExecution:
            - labelSelector:
                matchLabels:
                  app: api-gateway
              topologyKey: kubernetes.io/hostname
      containers:
        - name: gateway
          resources:
            requests: { cpu: "1", memory: "1Gi" }
            limits: { cpu: "2", memory: "2Gi" }
          readinessProbe:
            httpGet: { path: /health, port: 8080 }
            initialDelaySeconds: 5
          livenessProbe:
            httpGet: { path: /health, port: 8080 }
            initialDelaySeconds: 15

#8. 生产最佳实践

#8.1 性能优化

  • 启用 HTTP/2 连接复用
  • gRPC 连接池预热
  • 响应压缩(gzip/brotli)
  • 静态路由表缓存

#8.2 安全建议

  • 始终启用 TLS(最低 TLS 1.2)
  • 设置请求体大小限制(默认 10MB)
  • 启用请求日志(但脱敏敏感信息)
  • 定期轮换 JWT 签名密钥

#8.3 限流调优

  • 从保守的限流阈值开始,根据监控数据逐步调整
  • 为关键 API 设置独立的限流策略
  • 对内部服务调用使用更宽松的限流
  • 监控限流拒绝率,过高则需扩容

#9. 总结

coomia-dip API 网关作为平台的统一入口,实现了请求全生命周期的管理。关键设计亮点:

  1. 双协议支持:REST + gRPC 统一入口,REST-to-gRPC 自动转码
  2. 多维限流:全局/用户/IP/API 四维令牌桶,Redis 分布式实现
  3. 声明式路由:配置驱动的路由规则,支持动态服务发现
  4. gRPC-Web:前端应用直接调用 gRPC 服务
  5. 安全加固:JWT/OAuth2/API Key 多策略认证 + CORS + TLS

本篇是 S6 平台工程系列的最后一篇,完整覆盖了从权限控制、数据安全、SDK 开发到运维部署的全部平台工程能力。