找回密码
立即注册
搜索
发回帖 发新帖

6415

积分

0

好友

810

主题
发表于 2 小时前 | 查看: 5| 回复: 0

面向模型服务、后端与平台工程师。本文以 GPU 商品排序服务 ranker 为例,讲解如何用不可变制品、模型注册中心、预热就绪、Stable/Canary 容器池和流量加权完成无感更新。进程内热加载是可选优化,不是可靠发布的前提。示例中的流量、容量与阈值均为设计案例,须以本团队压测和 SLO 替换。

一次模型更新不该等同于"重启线上服务"。如果把旧模型进程停掉、再等待新进程下载权重并初始化 GPU,流量会落到尚未准备好的实例上,表现为 5xx、超时、连接中断或 P99 突刺。

更可靠的目标是:新版本在不接流量时完成下载、校验、加载和 dry run;只有通过就绪检查的实例才进入 Canary;流量由网格逐步转移;旧版本保留到观察期与请求排空完成。 这就是本文所说的"热加载与无感更新"。

1. 先看真实问题:为什么重启会让请求失败

假设 ranker 为商品列表请求做一次 GPU 前向计算:输入包含 80 个数值特征与上下文,输出 Top-N 排序分数。生产环境有 8 个 Pod,总流量约 600 QPS;模型权重约 7 GiB;冷启动、加载与代表性预热约 90 秒。模型、归一化参数与标签映射每周更新数次。

朴素发布通常这样发生:

停止 v41 Pod → 创建 v42 Pod → 下载权重 → 加载 GPU → 开端口 → 接收流量

问题在于"HTTP 端口已开"并不意味着模型可服务。若 readiness 只检查端口,v42 仍在加载时就会收到请求;若直接重启全部副本,稳定版本也不存在,无法快速回切。

本文的可验收目标不是"完全不重启",而是:

目标 可验证证据
未预热的 v42 不接收业务请求 actual_revision、/readyz、网格端点记录
v41 请求不被中途切换 路由排空记录;若使用进程内热加载则有代际测试
v42 失败时 v41 继续可用 候选失败演练、Stable 端点持续 ready
可以回滚 高 revision 指回 v41,或立即路由回切 Stable
每个结果可追溯 artifact、revision、runtime build 写入日志/响应元数据

2. 总体架构:先准备,再切流

这套方案的主角不是 Python 全局变量,而是发布控制面和两个可独立承载流量的服务池。

组件职责应明确:

组件 职责 不应承担的职责
发布流水线 产出不可变制品、生成 manifest、完成离线评估 直接修改线上 Pod 内模型
模型注册中心 保存制品元数据、审批状态、目标版本与审计 代理大文件推理流量
发布控制器 设置 Stable/Canary 目标、聚合状态、推进或回滚 在请求路径中加载模型
推理运行时 下载校验、加载、dry run、上报 actual、提供推理 擅自决定全量发布
服务网格 依照已批准权重路由请求 判断模型质量

这里的"动态加载"发生在每个推理运行时中。运行时只有在候选验证成功后才把自己标为可服务;流量切换永远由控制面和路由层决定。

3. 模型注册中心:制品身份与发布意图必须分开

3.1 制品是完整、不可变的运行单元

ranker-v42 不能只是一份权重。它至少携带权重、前后处理资源、输入契约和兼容性约束:

{
  "model": "ranker",
  "artifact_id": "ranker-v42",
  "runtime": "pytorch-state-dict",
  "feature_schema": "ranker-features-v3",
  "precision": "fp16",
  "files": [
    {"path": "weights.pt", "sha256": "<64-hex>", "bytes": 7516192768},
    {"path": "normalization.json", "sha256": "<64-hex>", "bytes": 8192},
    {"path": "label_mapping.json", "sha256": "<64-hex>", "bytes": 4096}
  ]
}

下载器写入唯一临时目录,逐文件检查大小和 SHA-256,全部通过后才将目录发布到本机内容寻址缓存。临时目录与最终目录应在同一文件系统;要求崩溃恢复时,需定义完成标记、fsync 与启动清扫策略。

摘要只证明内容与可信 manifest 一致,不能证明来源可信。manifest 必须来自经过认证和授权的注册中心;下载器还应拒绝路径穿越、未知文件、文件数超限和压缩解压膨胀。

3.2 回滚依赖 revision,不依赖"版本号变小"

artifact_id 表示不可变内容,revision 表示单调递增的发布意图:

revision=100  artifact=ranker-v41   Stable
revision=101  artifact=ranker-v42   Canary
revision=102  artifact=ranker-v41   回滚或重新指向已验证制品

若 v42 下载很慢,revision 102 已经发布,慢任务在提交前必须看到新目标并丢弃自己。相同 revision 指向不同 artifact 是控制面一致性错误,应拒绝并报警。

3.3 PostgreSQL 元数据与 CAS 发布

下面的骨架将 Stable 和 Canary 目标分开,避免它们共同订阅含糊的 latest:

CREATE TABLE model_artifact (
    model_name TEXT NOT NULL,
    artifact_id TEXT NOT NULL,
    manifest_uri TEXT NOT NULL,
    manifest_sha256 CHAR(64) NOT NULL,
    state TEXT NOT NULL CHECK (state IN ('pending', 'ready', 'rejected')),
    created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
    PRIMARY KEY (model_name, artifact_id)
);

CREATE TABLE deployment_target (
    model_name TEXT NOT NULL,
    environment TEXT NOT NULL,
    cohort TEXT NOT NULL CHECK (cohort IN ('stable', 'canary')),
    revision BIGINT NOT NULL CHECK (revision > 0),
    artifact_id TEXT NOT NULL,
    updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
    PRIMARY KEY (model_name, environment, cohort),
    FOREIGN KEY (model_name, artifact_id)
        REFERENCES model_artifact (model_name, artifact_id)
);

-- 在同一事务中先校验 artifact 为 ready,并插入审计记录。
UPDATE deployment_target
SET artifact_id = :artifact_id,
    revision = revision + 1,
    updated_at = now()
WHERE model_name = :model_name
  AND environment = :environment
  AND cohort = :cohort
  AND revision = :expected_revision;

调用方必须检查更新行数。零行说明并发发布冲突或目标不存在,应重新读取再决策,不能盲写。审计记录至少包括操作者、原因、旧新 artifact、旧新 revision 与审批信息。

4. 新模型如何安全进入 Canary:下载、加载、dry run、就绪

运行时维护两个状态:控制面期待的 desired 和它已成功加载的 actual。新目标到来后,执行以下流程。

候选必须完成以下检查,才能上报 actual_revision:

  1. 校验 manifest 内所有文件及缓存完成标记。
  2. 校验 feature_schema、精度、模型格式以及允许的 PyTorch/CUDA/驱动组合。
  3. 用服务内置的已知架构构造模型;绝不从制品执行任意 Python 代码。
  4. 严格加载 state dict,检查关键层名与参数形状。
  5. 用生产真实 shape 的代表性请求运行前处理、前向、后处理。
  6. 验证输出 shape、dtype、有限值及业务不变量,并等待 CUDA 预热完成。
  7. 检查加载峰值未超过 CPU、磁盘和 GPU 容量预算。

以下是运行时协调循环的实现骨架。具体下载器、注册中心客户端和指标实现应按团队 SDK 接入,但状态转换与失败处理不应省略:

from dataclasses import dataclass

@dataclass(frozen=True)
class DeploymentTarget:
    revision: int
    artifact_id: str
    manifest_uri: str
    manifest_sha256: str
    feature_schema: str

def reconcile(runtime, target: DeploymentTarget) -> None:
    """由单飞后台 worker 调用;不在 HTTP 请求处理线程中执行。"""
    if target.revision < runtime.desired_revision:
        return  # 迟到的控制面通知
    runtime.remember_desired(target)

    if runtime.actual_revision == target.revision:
        return

    try:
        artifact_dir = runtime.cache.fetch_verified(
            manifest_uri=target.manifest_uri,
            manifest_sha256=target.manifest_sha256,
        )
        candidate = runtime.loader.prepare(
            artifact_dir=artifact_dir,
            expected_schema=target.feature_schema,
        )
        runtime.loader.dry_run(candidate, runtime.representative_requests)

        # 安装前再读一次 desired,防止 v42 覆盖已经发布的回滚。
        if runtime.current_desired_revision() != target.revision:
            runtime.loader.dispose(candidate)
            return

        runtime.install(candidate, target)
        runtime.report_actual(target, status="ready")
    except Exception as exc:
        # 保留现有 Serving 模型;候选失败绝不把健康实例置为 NotReady。
        runtime.record_candidate_failure(target, exc)
        runtime.report_actual(target, status="failed")

候选加载应串行或单飞,并设置有界队列。若 target 不断变化,只保留最新尚未执行的任务;失败使用有上限的退避和抖动,避免每轮轮询都再次下载。

4.1 PyTorch dry run:输入契约比"随机跑一下"更重要

def prepare_ranker_candidate(weights_path, device, representative_requests):
    model = build_ranker_v3()  # 架构由服务代码拥有
    state = torch.load(weights_path, map_location="cpu", weights_only=True)
    model.load_state_dict(state, strict=True)
    model.to(device).eval()

    with torch.inference_mode():
        for request in representative_requests:
            inputs = move_to_device(request.model_inputs, device)
            output = model(**inputs)
            validate_ranker_output(output, request.expected_contract)

    if device.type == "cuda":
        torch.cuda.synchronize(device)  # 候选预热完成条件,不放在每个请求路径
    return model

eval() 与 inference_mode() 职责不同。代表性请求需覆盖实际 batch、shape、mask、稀疏/缺失特征和后处理路径;随机 Tensor 不能证明生产路径完成初始化。

5. Kubernetes 部署与 GPU 容量

探针必须区分首次启动、可接流量和进程存活:

端点 返回 200 的条件 返回 503 的条件
/startupz 首个可信模型已准备完成 尚未拥有可用模型
/readyz 有 Serving 模型,且不在排空 无可服务模型或正在终止
/livez 进程仍能安全运行 不可恢复的关键运行错误

候选 v42 失败而 v41 健康时,/readyz 应继续为 200;发布控制器根据 actual 判断 v42 是否可进入 Canary,而不是强迫稳定实例失去服务能力。

apiVersion: apps/v1
kind: Deployment
metadata:
  name: ranker-canary
spec:
  replicas: 2
  selector:
    matchLabels: {app: ranker, cohort: canary}
  template:
    metadata:
      labels: {app: ranker, cohort: canary}
    spec:
      nodeSelector: {workload: inference-gpu}
      terminationGracePeriodSeconds: 120
      containers:
      - name: ranker
        image: registry.example/ranker:runtime-20261008
        ports: [{containerPort: 8080}]
        startupProbe:
          httpGet: {path: /startupz, port: 8080}
          periodSeconds: 5
          failureThreshold: 120
        readinessProbe:
          httpGet: {path: /readyz, port: 8080}
          periodSeconds: 5
          failureThreshold: 2
        livenessProbe:
          httpGet: {path: /livez, port: 8080}
          periodSeconds: 15
          failureThreshold: 3
        resources:
          requests: {cpu: "2", memory: 8Gi}
          limits: {memory: 12Gi, nvidia.com/gpu: 1}

这不是可直接复制的容量值。发布前先算并实测:

GPU peak = weights + live inference + load temporary
         + warmup workspace + runtime + fragmentation margin

例如 24 GiB 卡上,7 GiB 权重、4 GiB 在途与运行时、3 GiB 加载/预热临时量、3 GiB 余量已约 17 GiB。若进程内同时装载旧 7 GiB 和新 8 GiB,则可能达到 25 GiB,不能双缓冲。torch.cuda.empty_cache() 只能释放分配器未使用的缓存块,不会释放仍被引用的张量。

GPU 推理服务的弹性伸缩不宜只看 CPU。可结合排队长度、并发数、P95/P99、GPU 利用率和显存水位扩缩容。HPA 负责副本数,节点扩缩容负责可调度 GPU;模型下载与 90 秒预热意味着扩容生效有延迟,Canary 发布也必须预留额外 GPU 容量。

6. 灰度发布:从 0% 到 100%,每一步都有停止条件

Canary 准备完成不代表立刻全量。Stable 和 Canary 保持独立目标:

Stable: revision=100, artifact=ranker-v41
Canary: revision=101, artifact=ranker-v42

在至少一个 Canary 实例上报 actual=(101, ranker-v42)、ready 容量足够后,才开始路由。以下 Istio 配置 表达 95/5 的一个阶段:

apiVersion: networking.istio.io/v1
kind: DestinationRule
metadata:
  name: ranker
spec:
  host: ranker
  subsets:
  - name: stable
    labels: {cohort: stable}
  - name: canary
    labels: {cohort: canary}
---
apiVersion: networking.istio.io/v1
kind: VirtualService
metadata:
  name: ranker
spec:
  hosts: [ranker]
  http:
  - route:
    - destination: {host: ranker, subset: stable}
      weight: 95
    - destination: {host: ranker, subset: canary}
      weight: 5
    timeout: 2s

一个可操作的推进表如下;时间与阈值来自团队 SLO,不应照搬:

阶段 Canary 流量 推进前提 停止或回滚条件
预热 0% Canary 全部校验、dry run、ready 任一候选校验/预热失败
小流量 5% 容量、探针、actual 一致 错误率、P99、显存或队列异常
扩大 25% / 50% 指标正常且样本量足够 技术或业务指标恶化
全量 100% 观察窗口满足门槛 立即路由回切 Stable
回收 100% 新版本 旧请求排空、观察期结束 仍保留旧制品的重新加载能力

5% 请求不必然代表 5% GPU 消耗:输入长度、batch、租户和缓存命中都会改变资源比例。用户级 A/B 还需要稳定分桶,不能只依赖随机权重路由。

回滚有两条路径:旧 Stable 仍就绪时,先把新请求的路由权重切回 Stable;旧池已回收时,发布更高 revision 指向旧 artifact,重新预热后再切流。二者都不能承诺无条件"秒级回滚"。

7. 多实例收敛、排空与失败处理

控制面保存 desired,数据面上报 actual。实例至少报告实例 ID、cohort、desired revision、actual revision、artifact、准备状态、最后错误、心跳与可服务容量。

控制面断连时:已验证的模型继续服务并告警;不接受未从可信注册中心获得的新制品;新 Pod 是否允许从可信本地缓存启动由恢复策略预先定义。重连后以权威 desired 重新收敛;同 revision 却 artifact 不同是严重一致性错误,不能自动猜测正确方。

缩容、回收或收到 SIGTERM 时,服务应先拒绝新准入、等待端点传播、排空在途请求,再释放资源。terminationGracePeriodSeconds 必须覆盖传播和排空预算;preStop 若存在,也消耗同一宽限期。

对普通单次前向推理,路由摘流后等待请求结束即可。对流式 LLM,旧会话必须持续绑定旧模型;KV Cache 与权重、Adapter、位置编码和运行时布局相关,不能因 Tokenizer 相同就跨模型复用。若团队同时在评估外部模型 API 做对照实验,RouteFast.ai 这类多模型中转服务也可以帮助降低切换与对比成本。

8. 高级优化:同一 Pod 内的代际热加载

跨 Pod 灰度已能实现服务不下线。只有当模型更新频率高、冷启动确实成为瓶颈、GPU 有新旧模型双缓冲余量,且团队已具备完善观测时,才需要在同一进程内切换模型。

此时的切换单位必须是完整代际:模型、artifact、revision、特征 schema、归一化参数、标签映射及 runtime build 一起借出。请求入场时增加该代际引用;发布新代际时退休旧代际;旧代际引用归零且真实 GPU 工作完成后再异步释放。

核心规则是:借用、读取版本和增加引用必须在同一把短锁内;锁不覆盖下载和推理;清理在锁外由可观测 worker 执行;CUDA 的非默认 stream、异步 D2H 拷贝或后台后处理要以相应 CUDA event 作为归还条件,而不是以 HTTP 返回、协程取消或 Python 函数返回作为归还条件。

以下是可直接运行的 Python 3.10+ 示例。它不依赖 PyTorch,专门验证最容易被实现错误破坏的并发语义:旧请求固定在 v41、新请求读取 v42、迟到候选不能覆盖回滚、清理只会在锁外排队一次。保存为 ranker_generation_demo.py 后执行 python3 ranker_generation_demo.py。

from __future__ import annotations

from contextlib import contextmanager
from dataclasses import dataclass
from queue import SimpleQueue
from threading import Barrier, Event, Lock, Thread
from typing import Iterator, Optional

class ToyRanker:
    def __init__(self, bias: int) -> None:
        self.bias = bias
        self.closed = False

    def score(self, value: int) -> int:
        if self.closed:
            raise RuntimeError("ranker already closed")
        return value + self.bias

    def close(self) -> None:
        self.closed = True

@dataclass(eq=False)
class Generation:
    revision: int
    artifact_id: str
    model: ToyRanker
    refs: int = 0
    retired: bool = False
    cleanup_queued: bool = False

class GenerationHolder:
    def __init__(self) -> None:
        self._lock = Lock()
        self._current: Optional[Generation] = None
        self._desired = (0, "")
        self._cleanup_queue: SimpleQueue[Generation] = SimpleQueue()

    def set_desired(self, revision: int, artifact_id: str) -> None:
        with self._lock:
            known_revision, known_artifact = self._desired
            if revision < known_revision:
                return
            if revision == known_revision:
                if artifact_id != known_artifact:
                    raise ValueError("same revision has different artifact")
                return
            self._desired = (revision, artifact_id)

    def _queue_cleanup_locked(self, generation: Generation) -> None:
        if not generation.cleanup_queued:
            generation.cleanup_queued = True
            self._cleanup_queue.put(generation)

    @contextmanager
    def borrow(self) -> Iterator[Generation]:
        with self._lock:
            generation = self._current
            if generation is None:
                raise RuntimeError("ranker is not ready")
            generation.refs += 1
        try:
            yield generation
        finally:
            with self._lock:
                generation.refs -= 1
                if generation.refs < 0:
                    raise AssertionError("unbalanced generation release")
                if generation.retired and generation.refs == 0:
                    self._queue_cleanup_locked(generation)

    def install(self, candidate: Generation) -> bool:
        # 调用方独占 candidate,且已完成真实模型加载和 dry run。
        with self._lock:
            current = self._current
            accepted = (
                (candidate.revision, candidate.artifact_id) == self._desired
                and (current is None or candidate.revision > current.revision)
            )
            if not accepted:
                self._queue_cleanup_locked(candidate)
                return False

            self._current = candidate
            if current is not None:
                current.retired = True
                if current.refs == 0:
                    self._queue_cleanup_locked(current)
            return True

    def drain_cleanup(self) -> None:
        # 生产中替换为可观测、可重试的后台 worker;异常不能回传给请求。
        while not self._cleanup_queue.empty():
            generation = self._cleanup_queue.get()
            try:
                generation.model.close()
            except Exception as exc:  # noqa: BLE001 - 演示异常处理点
                print(f"cleanup failed for {generation.artifact_id}: {exc}")

def main() -> None:
    holder = GenerationHolder()
    v41 = Generation(100, "ranker-v41", ToyRanker(41))
    holder.set_desired(100, "ranker-v41")
    assert holder.install(v41)

    # 旧请求已借到 v41;主线程发布 v42;旧请求最后才被允许结束。
    borrowed = Barrier(2)
    allow_old_request_to_finish = Event()
    old_result: list[int] = []

    def old_request() -> None:
        with holder.borrow() as generation:
            assert generation.artifact_id == "ranker-v41"
            borrowed.wait()
            allow_old_request_to_finish.wait(timeout=2)
            old_result.append(generation.model.score(1))

    thread = Thread(target=old_request)
    thread.start()
    borrowed.wait(timeout=2)

    holder.set_desired(101, "ranker-v42")
    v42 = Generation(101, "ranker-v42", ToyRanker(42))
    assert holder.install(v42)
    with holder.borrow() as fresh:
        assert (fresh.artifact_id, fresh.model.score(1)) == ("ranker-v42", 43)
    assert not v41.model.closed

    allow_old_request_to_finish.set()
    thread.join(timeout=2)
    assert not thread.is_alive()
    assert old_result == [42]
    holder.drain_cleanup()
    assert v41.model.closed

    # 回滚使用更高 revision;迟到的 v42 候选不会覆盖新的目标。
    holder.set_desired(102, "ranker-v41")
    stale_v42 = Generation(101, "ranker-v42", ToyRanker(42))
    assert not holder.install(stale_v42)
    v41_rollback = Generation(102, "ranker-v41", ToyRanker(41))
    assert holder.install(v41_rollback)
    holder.drain_cleanup()
    assert stale_v42.model.closed
    with holder.borrow() as current:
        assert (current.revision, current.artifact_id) == (102, "ranker-v41")

    print("PASS: pinning, concurrent switch, stale rejection, rollback")

if __name__ == "__main__":
    main()

真实项目可按以下边界拆分;这样正文中的注册中心、制品缓存、运行时和 HTTP 服务不会混在一个模块中:

ranker-service/
├── app.py                 # /predict、/startupz、/readyz、/livez 与优雅终止
├── registry_client.py     # desired target 拉取、actual 状态上报
├── artifact_cache.py      # 临时目录下载、manifest 校验、原子缓存发布
├── model_loader.py        # PyTorch 构建、strict load、dry run 与 CUDA event
├── generation_holder.py   # 本文的 GenerationHolder
├── rollout_metrics.py     # Prometheus 指标和审计字段
└── tests/
    ├── test_artifact_cache.py
    ├── test_reconcile.py
    └── test_generation_switch.py

最少应自动验证:旧请求跨越切换仍用 v41;新请求只见 v42;v42 慢加载不能覆盖 revision 102 的回滚;多个请求同时结束只会安排一次清理;清理失败不会改变已经成功的响应。这一层优化不替代 Stable/Canary 路由,而是为单 Pod 减少冷启动成本。

9. 上线验收:用证据证明"无感"

发布前先写验收门槛。最低可观测性包括:Stable/Canary 的 QPS、错误率、P95/P99;排队与拒绝数;下载/加载/dry run/提交耗时;desired/actual 不一致;GPU allocated/reserved 与设备层占用;每个 artifact、revision 和 runtime build;输出质量或业务代理指标。

故障注入 期望结果
制品截断、摘要或 schema 不符 Canary 不 ready,Stable 持续服务
v42 慢加载期间发出回滚 迟到候选不能覆盖更高 revision
dry run 通过但灰度 P99 恶化 停止推进,路由回切 Stable
Canary GPU OOM 实例摘流,Stable 容量保持服务
控制面断连或重复通知 已验证模型继续服务,恢复后重新收敛
SIGTERM 与切流同时发生 停止新准入、排空后关闭,不重新接流

真正交付的不是一个"热加载开关",而是一条证据链:制品可信、候选已预热、未就绪实例不接流量、流量切换可控、旧版本可以承接回滚、每次结果都可追溯。做到这些,模型服务就能在不下线的前提下平滑过渡到新版本。


参考资料

  1. PyTorch CUDA semantics:异步执行、stream 与同步语义。
  2. PyTorch empty_cache:缓存分配器行为边界。
  3. PyTorch inference_mode:线程局部推理上下文。
  4. Kubernetes Probes:startup、readiness、liveness 的职责。
  5. Istio Traffic Shifting:按权重迁移流量。



上一篇:OpenAI新论文解读:Landau-Siegel零点问题的证明思路与意义
下一篇:打工十年才悟透:财务自由靠的不是辞职创业,而是先上班攒这三样
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-10-10 05:14 , Processed in 0.085095 second(s), 41 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表