全部术语

GLOSSARY ENTRY

服务可靠性、Backpressure 与 Observability

  • Serving Reliability
  • Backpressure
  • Observability
  • Metrics
  • Distributed Tracing

用有界流控、幂等终止和资源所有权不变量限制慢客户端与取消竞态的影响,再用低基数指标、请求级 trace 和结构化日志定位 SLO、排队、GPU 与清理问题。

它保护控制平面,不改变模型数学

服务可靠性位于 Request Lifecycle、scheduler、transport 与 worker 之间。它不改变 attention、sampling 或 token 序列的理论分布,而是保证:慢客户端、deadline、取消、worker error 或部分组件重试发生时,请求不会无限占用 host buffer、KV blocks、batch slots 和后台计算。

可检查的终止不变量是:

terminal(r){stop,length,cancel,timeout,error}terminal(r)\in\{stop,length,cancel,timeout,error\} cleanup_count(r)=1,KVr=0,slotsr=0cleanup\_count(r)=1, \qquad KV_r=0,\quad slots_r=0

“只有一个终态”不要求只有一个终止信号到达;deadline、client disconnect 和 backend error 可能竞态。它要求第一个成功提交终态的 owner 赢得 compare-and-set/actor mailbox,后续信号观察到已终止并成为 no-op。若共享 Prefix Cache blocks,KV_r=0 指请求引用被释放,不是物理 block 无条件删除。

Backpressure 要逐层传播

令 tick tt 的待发送队列长度为 QtQ_t,本轮生产 gtg_t 个 chunks、transport 消费 dtd_t 个,容量上限为 CC

Qt+1=min(C,max(0,Qtdt)+gt)Q_{t+1}=\min\left(C,\max(0,Q_t-d_t)+g_t\right)

但只有在 Q < C 时仍允许 gt>0g_t>0,这个式子才表示真实 backpressure;若生产者继续生成、只是覆盖或丢弃,内存有界却浪费 GPU/KV 工作。无界队列则近似:

Qt+1=max(0,Qtdt)+gtQ_{t+1}=\max(0,Q_t-d_t)+g_t

当平均生产率 λg>λd\lambda_g>\lambda_dQQ 随时间增长,直到请求结束、进程内存耗尽或外层 timeout。

一条完整反馈链通常是:

socket / HTTP2 / SSE flow-control window exhausted

bounded per-request async send queue reaches high watermark

streamer stops accepting more committed chunks

scheduler temporarily stops advancing this request

deadline / slow-consumer policy may finalize and release KV

只在 socket writer 阻塞,却让 model worker 继续产生 token,不能限制 generated-but-unsent 状态。反过来,让一个慢客户端阻塞整个 scheduler thread,会把局部背压扩散成所有请求的 head-of-line blocking。生产实现常把 transport、detokenization 与 GPU scheduling 隔离,并用 per-request credits/high-low watermarks 传递状态。

四种 token 计数不能混成一个

计数所有权边界终止时可能怎样
proposeddraft/MTP 产生,target 尚未接受可被 verifier 拒绝,不得发布
committedtarget 已接受并进入请求逻辑序列/KV应成为计费/finish 语义的明确候选
generated/queued已转为待发送 chunk,仍在 host buffer慢客户端取消时可能未发送而丢弃
sent/acknowledged已交给 transport;是否被客户端应用层读取通常未知不能靠回滚 KV 撤回

因此 tokens/s 必须说明分子。投机解码用 proposed tokens 会虚高,streaming 服务用 generated 代替 sent 会掩盖 backpressure waste。

一个慢客户端 worked trace

生产 token 为 A..G,每 tick 最多生成一个;drain trace 为 [1,0,0,0,1,0,1,0,0],tick 8 到达 overall deadline。这里 transport 在生成前 drain,因此 tick 0 还没有内容可发送:

tickbounded capacity=3:sentpendingscheduler 动作
0""A生成 A
1""AB生成 B
2""ABC生成 C,达到上限
3""ABC暂停生成
4ABCDdrain A,再生成 D
5ABCD暂停生成
6ABCDEdrain B,再生成 E
7ABCDE暂停生成
8ABCDEdeadline;丢弃未发送并清理

无界路径在同一时刻已生成 ABCDEFG,仍只发送 AB,终止时丢弃 CDEFG;有界路径只生成 ABCDE,未发送上界为 3。这个例子没有宣称暂停一定提高吞吐:若 client 很快恢复,过小 buffer 可能增加 scheduler wakeups、破坏 batch efficiency 或引入 TPOT 抖动。

可运行的流控、终止竞态与遥测实验

实验问题:同一慢 consumer 下 bounded queue 是否限制 buffer 和无效生成?overall deadline 与 client cancel 同 tick 到达时,资源是否只回收一次?同一请求摘要能否同时产生 SLO goodput、低基数 metric labels 与含 request id 的 traces?

import platform,statistics
from collections import Counter
from dataclasses import dataclass,field

@dataclass
class ResourceOwner:
    request_id:str; kv_blocks:int; slots:int=1
    terminal_reason:str|None=None; cleanup_count:int=0
    finalize_attempts:list[str]=field(default_factory=list)
    def finalize(self,reason):
        self.finalize_attempts.append(reason)
        if self.terminal_reason is not None: return False
        self.terminal_reason=reason; self.kv_blocks=0; self.slots=0
        self.cleanup_count+=1; return True

def simulate_stream(capacity,deadline_tick):
    tokens=list("ABCDEFG"); drain_per_tick=[1,0,0,0,1,0,1,0,0]
    owner=ResourceOwner("req-42",kv_blocks=4)
    buffer=[]; sent=[]; generated=[]; paused=[]; trace=[]
    for tick,drain in enumerate(drain_per_tick):
        if tick>=deadline_tick:
            assert owner.finalize("overall_deadline")
            assert not owner.finalize("client_cancel")
            trace.append((tick,"terminal","".join(buffer))); break
        for _ in range(min(drain,len(buffer))): sent.append(buffer.pop(0))
        if len(generated)<len(tokens):
            if capacity is None or len(buffer)<capacity:
                token=tokens[len(generated)]; generated.append(token); buffer.append(token)
            else: paused.append(tick)
        trace.append((tick,"".join(sent),"".join(buffer)))
    assert owner.cleanup_count==1 and owner.kv_blocks==owner.slots==0
    if capacity is not None: assert max(len(state[2]) for state in trace)<=capacity
    return {"sent":sent,"generated":generated,"dropped":buffer,
            "paused_ticks":paused,"trace":trace,"owner":owner}

@dataclass(frozen=True)
class RequestSummary:
    request_id:str; finish_reason:str; arrival:float
    first_token:float|None; token_times:tuple[float,...]

def aggregate_telemetry(requests,window_seconds,ttft_slo,tpot_slo):
    finishes=Counter(r.finish_reason for r in requests)
    ttft_samples=[]; tpot_samples=[]; good=0; traces=[]
    for request in requests:
        traces.append({"request.id":request.request_id,
                       "finish.reason":request.finish_reason})
        if request.first_token is None: continue
        ttft=request.first_token-request.arrival
        gaps=[b-a for a,b in zip(request.token_times,request.token_times[1:])]
        tpot=statistics.mean(gaps) if gaps else 0.
        ttft_samples.append(ttft); tpot_samples.append(tpot)
        if request.finish_reason=="stop" and ttft<=ttft_slo and tpot<=tpot_slo: good+=1
    metric_labels=[{"finish_reason":reason} for reason in sorted(finishes)]
    assert all("request_id" not in labels for labels in metric_labels)
    return {"finish_total":dict(finishes),"ttft_samples":ttft_samples,
            "tpot_samples":tpot_samples,"goodput_rps":good/window_seconds,
            "metric_label_sets":metric_labels,"trace_attributes":traces}

def main():
    unbounded=simulate_stream(None,8); bounded=simulate_stream(3,8)
    assert bounded["sent"]==unbounded["sent"]==list("AB")
    assert len(bounded["generated"])<len(unbounded["generated"])
    assert len(bounded["dropped"])<=3<len(unbounded["dropped"])
    requests=[RequestSummary("req-a","stop",0.,2.,(2.,3.,4.)),
              RequestSummary("req-b","stop",0.,5.,(5.,8.)),
              RequestSummary("req-c","client_cancel",1.,None,())]
    telemetry=aggregate_telemetry(requests,10.,ttft_slo=3.,tpot_slo=2.)
    assert telemetry["finish_total"]=={"stop":2,"client_cancel":1}
    assert telemetry["goodput_rps"]==.1
    print(f"python={platform.python_version()} implementation={platform.python_implementation()}")
    print("stream_tokens=ABCDEFG drain_per_tick=[1,0,0,0,1,0,1,0,0] deadline_tick=8")
    for name,result in (("unbounded",unbounded),("bounded_capacity_3",bounded)):
        print(f"{name}: sent={''.join(result['sent'])} generated={''.join(result['generated'])} "
              f"dropped={''.join(result['dropped'])} paused_ticks={result['paused_ticks']}")
    owner=bounded["owner"]
    print(f"terminal_reason={owner.terminal_reason} finalize_attempts={owner.finalize_attempts} "
          f"cleanup_count={owner.cleanup_count} kv_blocks={owner.kv_blocks} slots={owner.slots}")
    print(f"bounded_trace={bounded['trace']}")
    print(f"finish_total={telemetry['finish_total']} ttft_samples={telemetry['ttft_samples']} "
          f"tpot_samples={telemetry['tpot_samples']} goodput_rps={telemetry['goodput_rps']:.3f}")
    print(f"metric_label_sets={telemetry['metric_label_sets']}")
    print(f"trace_attributes={telemetry['trace_attributes']}")

if __name__=="__main__": main()

完整文件位于 examples/reliability_observability.py。实际无界/容量 3 路径分别生成 ABCDEFG / ABCDE,终止时丢弃 CDEFG / CDE;bounded scheduler 在 ticks [3,5,7] 暂停。deadline 先赢得终态,随后的 cancel 返回 no-op,cleanup_count=1 且 KV/slot 均为 0。

三个请求在 10 秒窗口中产生 stop=2, client_cancel=1;TTFT samples 为 [2,5],TPOT samples 为 [1,3]。在 TTFT≤3s, TPOT≤2s 且只把正常完成计为 good 的口径下,goodput 为 0.1 request/s

Metrics、traces、logs 各回答什么

Signal适合回答典型字段不适合
Metrics“系统整体是否违反 SLO、容量是否逼近上限”queue depth、active KV tokens、TTFT/TPOT histogram、finish/reject/preempt counters为单个 request 保存任意细节
Distributed traces“req-42 时间花在哪,跨 gateway/scheduler/worker 的因果链是什么”request/span id、queue/prefill/decode/stream spans、batch/shape/worker attributes不采样地聚合全量长期趋势
Structured logs/events“发生了哪次状态跃迁或异常,携带什么离散上下文”timestamp、request hash、from/to state、reason、resource delta用文本搜索替代精确定义的 histogram/counter

Metrics labels 应低基数、稳定,例如 model、status/finish reason、priority class、worker pool;不要把 request id、prompt、adapter name 的任意用户值直接放进 labels,否则 time series 数量近似相乘:

NseriesjlabeljN_{series}\approx\prod_j |label_j|

request id 适合作为 trace attribute/exemplar 关联到个例。日志/trace 中的 prompt、输出、tenant id 还涉及 PII、商业数据与 prompt injection 内容,需 redaction、采样、访问控制和 retention policy;“为了排障全量记录文本”不是安全默认值。

Histogram 与 goodput 口径

TTFT、queue time、TPOT 应报告分布而非只报平均值。直方图 bucket 需覆盖真实 SLO 阈值,聚合时保留 count/sum/buckets;客户端侧 end-to-end TTFT 与 server model-side TTFT 不同,二者都应命名清楚。

若 SLO 为 TTFTx,TPOTyTTFT\le x,TPOT\le y,窗口 WW 内:

G=#{r:finishrSeligible,TTFTrx,TPOTry}WG=\frac{\#\{r:finish_r\in S_{eligible},TTFT_r\le x,TPOT_r\le y\}}{W}

SeligibleS_{eligible} 是否只含 stop/length,取消是否进入 attempted load 分母,必须固定。只从成功请求计算 latency 会产生 survivor bias:过载时最慢请求超时/取消后消失,成功样本 p99 甚至可能“改善”。至少联看 arrival/admission/reject/cancel/error totals 与 in-flight requests。

生产中真正难的边界

  • 异步竞态:socket disconnect、deadline timer、scheduler preemption 和 worker failure 分属不同线程/进程;教学对象的顺序调用要换成原子状态或单 owner 消息序列。
  • GPU 与 collective 安全点:cancel flag 不能抢占已 launch kernel;TP/EP ranks 必须保持 collective 顺序,再统一移除 request。
  • 资源图而非一个计数:KV pages、shared prefix refcounts、LoRA handles、CUDA Graph slots、stream buffers、trace contexts 可能有不同 owners;cleanup 要按 dependency/ownership 回收。
  • 进程崩溃:内存内 finally 不会在机器失联后执行;coordinator 需要 lease/heartbeat、orphan detection 与 worker restart/reconciliation。
  • 流控层错配:HTTP/2 window、async queue、SSE flush、proxy buffer 与客户端 SDK 各有缓存;只看应用 queue 长度不足以证明 bytes 已离开系统。
  • 可观测性成本:全量 per-token spans/logs 会增加 CPU、内存、网络和 tail latency;使用 request/phase spans、histograms、sampling/exemplars,并测量 instrumentation overhead。
  • 告警不是单指标阈值:TTFT 上升可能来自 queue/admission、prefill kernel、prefix miss 或 transport;把 latency、queue、KV usage、batch composition、GPU/collective 和 finish reasons 关联起来才可定位。

与推理服务主链路的连接

  • Request Lifecycle定义状态与 stop/cancel/timeout 语义;本页把慢 consumer、终态竞态、资源清理和观测证据展开成可验证不变量。
  • Scheduler Budget / Admission Control决定能否推进请求;bounded stream credits 可以成为额外 eligibility 条件,deadline/reject/preempt counters 则验证 policy 结果。
  • Continuous Batching优化 active batch;若它无视慢客户端,可能把 GPU work 转为 generated-but-unsent waste。
  • Benchmark、容量模型与 Profiler定义实验计时、内存和 goodput;本页定义在线 metrics/traces/logs 如何持续证明这些口径。
  • Prefill/Decode提供 TTFT/TPOT 阶段定义;观测时还必须区分 server-side model time、queue time 与 client-perceived stream time。

参考资料