LLM 推理网关生产级架构

LLM 推理网关生产级架构:动态模型路由、流量调度与灰度发布

本文面向构建生产级 LLM (Large Language Model) 推理平台的工程师,深入探讨推理网关 (Model Gateway) 的设计哲学与工程实现。

一、从单体到智能网关

早期 LLM 推理服务通常采用直接暴露模式:客户端发送请求到单一端点,后端绑定固定模型。随着多模型、多版本、多租户的演进,这种架构面临严峻挑战:

  • 模型碎片化:同时部署 7B/13B/70B 量化与非量化版本,需要统一入口。
  • 流量潮汐:高峰期某些模型 QPS 激增,需要灵活的弹性调度。
  • 版本灰度:新模型上线需要可控的流量切分与回滚策略。
  • 租户隔离:企业级场景需要基于 API Key 的配额管理与优先级路由。

推理网关作为客户端与推理引擎之间的统一接入层,承担了请求路由、协议适配、流量治理、可观测性等核心职责。本文将从架构设计、核心算法、生产实践三个维度展开。

二、核心架构设计

2.1 分层架构

采用分层解耦设计,自上而下依次为:

层级 职责
协议层 支持 OpenAI / Anthropic / gRPC 多协议归一化
路由层 基于 model 字段匹配目标模型池
治理层 限流、熔断、重试、超时、降级
调度层 权重分流、灰度比例、A/B 测试
适配层 转换 vLLM / TGI / Triton 等异构后端协议

2.2 请求生命周期


Client → [TLS Termination] → [AuthN/AuthZ] → [Route Match]
    → [Rate Limit Check] → [Queue & Weight Sharding]
    → [Backend Selection] → [Endpoint Proxy] → [Stream Response]

关键设计要点:每个阶段都是可插拔的中间件,支持运行时热更新配置。

三、动态路由策略

3.1 声明式路由规则

路由配置采用声明式 YAML,支持版本热加载:


routes:
  - name: default
    match:
      model: ["gpt-4", "gpt-4-turbo"]
    pool: gpt4-cluster
    strategy: weighted-round-robin
    
  - name: budget-tenant
    match:
      header:
        X-Tenant-Tier: "free"
    pool: shared-7b-pool
    strategy: least-requests
    rate_limit:
      rpm: 60
      tpm: 10000
    
  - name: canary-v2
    match:
      model: "llama-3-70b"
    pool: llama-cluster
    strategy: weighted-split
    weights:
      v1: 90
      v2: 10

3.2 权重分流的实现

加权随机分流是灰度发布的基础。与传统的整数权重(如 90:10)不同,生产环境需要支持平滑调整与一致性哈希:


use std::sync::atomic::{AtomicU32, Ordering};

pub struct WeightedRouter {
    backends: Vec<Backend>,
    cumulative_weights: Vec<f64>,
    total_weight: AtomicU32,
}

impl WeightedRouter {
    pub fn route(&self, request_id: &str) -> Option<&Backend> {
        // 一致性哈希将同一 session 固定到同一版本
        if let Some(v) = self.consistent_hash_route(request_id) {
            return Some(v);
        }
        
        // 否则按权重随机分流
        let total = self.total_weight.load(Ordering::Relaxed) as f64;
        let mut val = fastrand::f64() * total;
        
        for (i, weight) in self.cumulative_weights.iter().enumerate() {
            val -= weight;
            if val <= 0.0 {
                return Some(&self.backends[i]);
            }
        }
        self.backends.last()
    }
    
    pub fn adjust_weights(&self, weights: &[f64]) {
        // 原子更新,无锁切换
        let partial_sum: f64 = weights.iter().sum();
        let mut cumulative = 0.0;
        for (i, w) in weights.iter().enumerate() {
            cumulative += w / partial_sum;
            // 原子写入 cumulative_weights[i]
        }
        self.total_weight.store(1000, Ordering::Release);
    }
}

3.3 基于延迟的自适应路由

静态权重无法感知后端实时负载。生产级网关需要根据 P99 延迟动态调整后端权重:


pub struct AdaptiveRouter {
    inner: WeightedRouter,
    latency_tracker: Arc<LatencyTracker>,
    alpha: f64, // 指数移动平均系数
}

impl AdaptiveRouter {
    pub fn recompute_weights(&self) {
        let scores: Vec<f64> = self.inner.backends.iter().map(|b| {
            let p99 = self.latency_tracker.p99(&b.id);
            let success_rate = self.latency_tracker.success_rate(&b.id);
            let queue_depth = b.pending_requests() as f64;
            
            // 综合评分:高延迟、低成功率、大队列 → 低权重
            let score = (1.0 / (1.0 + p99 / 1000.0)) * success_rate;
            score / (1.0 + queue_depth * 0.01)
        }).collect();
        
        // Softmax 归一化
        let max_score = scores.iter().cloned().fold(0.0, f64::max);
        let exps: Vec<f64> = scores.iter().map(|s| ((s - max_score) / 100.0).exp()).collect();
        let sum_exp: f64 = exps.iter().sum();
        let weights: Vec<f64> = exps.iter().map(|e| e / sum_exp).collect();
        
        self.inner.adjust_weights(&weights);
    }
}

四、流量治理:限流与背压

4.1 多层次限流策略


// 基于令牌桶的全局限流
pub struct RateLimiter {
    // 全局 RPM 限制
    global_rpm: Arc<TokenBucket>,
    // 每模型 TPMin 限制(防长 prompt 耗尽配额)
    per_model_tpm: DashMap<String, TokenBucket>,
    // 每租户并发限制
    per_tenant_concurrency: DashMap<String, Semaphore>,
}

impl RateLimiter {
    pub async fn acquire(&self, ctx: &RequestContext) -> Result<(), RateLimitError> {
        // 检查全局 RPM
        self.global_rpm.acquire(1).await?;
        
        // 检查模型级 TPM(预估算 token)
        let estimated_tokens = ctx.prompt_tokens + ctx.max_tokens;
        self.per_model_tpm
            .entry(ctx.model.clone())
            .or_insert_with(|| TokenBucket::new(/* model tpmin */))
            .acquire(estimated_tokens as u32)
            .await?;
        
        // 检查租户并发
        let sem = self.per_tenant_concurrency
            .entry(ctx.tenant_id.clone())
            .or_insert_with(|| Semaphore::new(10));
            
        match tokio::time::timeout(
            Duration::from_millis(100),
            sem.acquire()
        ).await {
            Ok(permit) => {
                permit.forget(); // 由 drop guard 管理
                Ok(())
            }
            Err(_) => Err(RateLimitError::ConcurrencyLimit),
        }
    }
}

4.2 LLM 感知的背压机制

不同于普通 HTTP 服务,LLM 推理请求具有长尾特性(长文本生成可能持续数分钟)。传统连接超时会导致资源泄漏。

解决方案:采用分级超时 + 流式心跳


pub async fn proxy_with_backpressure(
    mut response: Response,
    stream_tx: tokio::sync::mpsc::Sender<Bytes>,
) -> Result<()> {
    let mut idle_timer = tokio::time::interval(Duration::from_secs(5));
    let deadline = Instant::now() + Duration::from_secs(120);
    
    loop {
        tokio::select! {
            chunk = response.chunk() => {
                match chunk {
                    Some(Ok(data)) => {
                        idle_timer.reset().await;  // 重置空闲计时器
                        stream_tx.send(data).await?;
                    }
                    Some(Err(e)) => return Err(e.into()),
                    None => return Ok(()),  // 正常结束
                }
            }
            _ = idle_timer.tick() => {
                // 5 秒内无数据,发送 SSE 心跳注释
                stream_tx.send(Bytes::from(": ping\n\n")).await?;
            }
            _ = tokio::time::sleep_until(deadline.into()) => {
                return Err(GatewayError::StreamTimeout);
            }
        }
    }
}

五、灰度发布与回滚

5.1 渐进式流量切分

生产环境灰度不是一步到位,而是渐进式推进。以下是自动化 Canary 流程:


traffic=5% → (观察 5min) → error_rate>1%? 回滚 : traffic=25%
→ (观察 10min) → p99>baseline*1.2? 回滚 : traffic=50%
→ (观察 30min) → 指标正常 → traffic=100%

实现一个状态机驱动的发布控制器:


#[derive(Debug, Clone, Copy, PartialEq)]
enum CanaryPhase {
    Initial,      // 5%
    Expanding,    // 25%
    Validating,   // 50%
    RollingOut,   // 100%
    RollingBack,
}

pub struct CanaryController {
    current_phase: AtomicU64,
    router: Arc<AdaptiveRouter>,
    metrics: Arc<MetricsCollector>,
}

impl CanaryController {
    pub async fn run(&self) -> Result<()> {
        let phases = vec![
            (CanaryPhase::Initial, 5, Duration::from_secs(300)),
            (CanaryPhase::Expanding, 25, Duration::from_secs(600)),
            (CanaryPhase::Validating, 50, Duration::from_secs(1800)),
            (CanaryPhase::RollingOut, 100, Duration::ZERO),
        ];
        
        for (phase, percentage, duration) in phases {
            self.router.set_canary_weight(percentage);
            
            tokio::time::sleep(duration).await;
            
            let report = self.metrics.evaluate_canary().await;
            if report.error_rate > 0.01 || report.p99_latency > report.baseline_p99 * 1.2 {
                self.router.set_canary_weight(0);
                tracing::error!(phase = ?phase, report = ?report, "Canary rolled back");
                return Err(CanaryError::ValidationFailed);
            }
            
            tracing::info!(phase = ?phase, "Canary phase passed");
        }
        Ok(())
    }
}

5.2 影子流量验证

在高风险场景下,可以利用影子流量 (Shadow Traffic) 在不影响线上流量的情况下验证新版本:


pub async fn shadow_dispatch(
    request: Request,
    primary: &Backend,
    shadow: &Backend,
) -> Result<Response> {
    // 主路正常转发
    let primary_fut = primary.handle(request.clone());
    
    // 影子路异步转发,结果仅用于对比
    let shadow_handle = tokio::spawn(async move {
        let response = shadow.handle(request).await;
        metrics::histogram!("shadow.latency", response.latency);
        metrics::counter!("shadow.errors", response.is_err() as u32);
        response  // 丢弃响应内容,只保留指标
    });
    
    // 等待主路返回结果
    match primary_fut.await {
        Ok(resp) => {
            // 异步等待影子路完成(不阻塞返回)
            let _ = shadow_handle.await;
            Ok(resp)
        }
        Err(e) => {
            tracing::error!("Primary backend failed, aborting shadow");
            shadow_handle.abort();
            Err(e)
        }
    }
}

六、生产级 Rust 实现

6.1 完整网关骨架


use axum::{
    routing::post,
    Router, extract::Extension,
    middleware::{self, Next},
};
use std::sync::Arc;

#[tokio::main]
async fn main() {
    // 初始化 tracing
    tracing_subscriber::fmt()
        .with_env_filter("model_gateway=debug,tower_http=debug")
        .init();
    
    // 构建应用状态
    let state = Arc::new(AppState::new(Config::from_env()).await);
    
    let app = Router::new()
        .route("/v1/chat/completions", post(handle_chat))
        .route("/v1/embeddings", post(handle_embed))
        .route("/health", axum::routing::get(health_check))
        .layer(middleware::from_fn_with_state(
            state.clone(),
            auth_middleware,
        ))
        .layer(middleware::from_fn_with_state(
            state.clone(),
            rate_limit_middleware,
        ))
        .layer(Extension(state.clone()));
    
    let listener = tokio::net::TcpListener::bind("0.0.0.0:8080").await.unwrap();
    axum::serve(
        listener,
        app.into_make_service_with_connect_info::<SocketAddr>(),
    )
    .tcp_nodelay(true)
    .await
    .unwrap();
}

6.2 可观测性三要素


pub struct Observability {
    request_counter: Counter<u64>,
    latency_histogram: Histogram<f64>,
    token_counter: Counter<u64>,
    active_gauge: Gauge<u64>,
}

impl Observability {
    pub fn record_request(&self, ctx: &RequestContext, duration: Duration, tokens: u32) {
        let labels = [
            ("model", ctx.model.as_str()),
            ("tenant", ctx.tenant_tier.as_str()),
            ("status", if ctx.success { "200" } else { "500" }),
        ];
        
        self.request_counter.inc_by(1, &labels);
        self.latency_histogram.observe(duration.as_secs_f64(), &labels);
        self.token_counter.inc_by(tokens as u64, &labels);
    }
}

// 暴露 Prometheus /metrics 端点
async fn metrics(Extension(obs): Extension<Arc<Observability>>) -> String {
    let encoder = prometheus::TextEncoder::new();
    let metric_families = prometheus::gather();
    encoder.encode_to_string(&metric_families).unwrap()
}

6.3 连接池与 HTTP/2 多路复用

LLM 推理网关与后院的连接管理至关重要。推荐使用 hyper 的连接池:


use hyper_util::rt::TokioExecutor;
use hyper::client::conn::http2::Builder;

pub struct BackendPool {
    clients: ArrayQueue<SendRequest<Body>>,
    endpoint: Uri,
}

impl BackendPool {
    pub async fn get(&self) -> Result<PooledConnection> {
        if let Some(conn) = self.clients.pop() {
            if conn.is_ready() {
                return Ok(PooledConnection::new(conn, self.clients.clone()));
            }
        }
        
        // 池中无可用连接,建立新连接(带 HTTP/2 多路复用)
        let conn = Builder::new(TokioExecutor::new())
            .initial_stream_window_size(1024 * 1024)  // 1MB 流窗口
            .initial_connection_window_size(4 * 1024 * 1024)
            .max_concurrent_streams(100)
            .handshake(self.endpoint.clone())
            .await?;
            
        Ok(PooledConnection::new(conn, self.clients.clone()))
    }
}

七、性能优化实践

7.1 零拷贝请求转发


// 利用 Bytes 的引用计数实现零拷贝转发
pub async fn zero_copy_proxy(
    request: Request,
    backend: &BackendPool,
) -> Result<Response> {
    let (parts, body) = request.into_parts();
    
    // 复用 Bytes,避免反序列化
    let body_bytes = axum::body::to_bytes(body, usize::MAX).await?;
    
    // 直接透传,不做 JSON 解析
    let forwarded = Request::from_parts(parts, Body::from(body_bytes.clone()));
    
    let mut conn = backend.get().await?;
    let response = conn.send_request(forwarded).await?;
    
    Ok(response)
}

7.2 Prompt Cache 友好设计

vLLM 等推理引擎支持前缀 KV Cache 复用。网关可通过路由关联最大化缓存命中:


pub fn cache_aware_route(&self, ctx: &RequestContext) -> Option<&Backend> {
    let prompt_fingerprint = blake3::hash(&ctx.system_prompt.as_bytes());
    let shard_idx = prompt_fingerprint.as_bytes()[0] as usize % self.backends.len();
    
    // 相同 system prompt 的请求路由到同一 backend
    // 最大化 RadixAttention 前缀树命中率
    Some(&self.backends[shard_idx])
}

八、总结与展望

LLM 推理网关是连接用户请求与推理引擎的关键中间层。本文阐述的核心要点可归纳为:

  1. 路由层声明化:将路由规则从代码中解耦,支持运行时热更新。
  2. 流量治理分层:全局 RPM → 模型级 TPM → 租户并发,层层递进。
  3. 自适应调度:结合权重、延迟、队列深度动态调整后端选择。
  4. 灰度自动化:状态机驱动 + 指标回判,实现无人值守发布。
  5. Rust 工程实践:利用所有权系统实现无锁高性能,tokio 异步生态保障吞吐。
  6. 随着推理引擎向 Disaggregated 架构(Prefill/Decode 分离)演进,网关的职责将进一步扩展——未来的推理网关需要理解请求的 Prefill 成本,智能调度到不同的推理阶段节点。但这已是下一个话题了。

    ---

    关键术语表: Model Gateway, Weighted Routing, Canary Release, Shadow Traffic, Rate Limiting, Token Bucket, Adaptive Routing, KV Cache Affinity, HTTP/2 Multiplexing, Backpressure

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部