它保护控制平面,不改变模型数学
服务可靠性位于 Request Lifecycle、scheduler、transport 与 worker 之间。它不改变 attention、sampling 或 token 序列的理论分布,而是保证:慢客户端、deadline、取消、worker error 或部分组件重试发生时,请求不会无限占用 host buffer、KV blocks、batch slots 和后台计算。
可检查的终止不变量是:
“只有一个终态”不要求只有一个终止信号到达;deadline、client disconnect 和 backend error 可能竞态。它要求第一个成功提交终态的 owner 赢得 compare-and-set/actor mailbox,后续信号观察到已终止并成为 no-op。若共享 Prefix Cache blocks,KV_r=0 指请求引用被释放,不是物理 block 无条件删除。
Backpressure 要逐层传播
令 tick 的待发送队列长度为 ,本轮生产 个 chunks、transport 消费 个,容量上限为 :
但只有在 Q < C 时仍允许 ,这个式子才表示真实 backpressure;若生产者继续生成、只是覆盖或丢弃,内存有界却浪费 GPU/KV 工作。无界队列则近似:
当平均生产率 , 随时间增长,直到请求结束、进程内存耗尽或外层 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 计数不能混成一个
| 计数 | 所有权边界 | 终止时可能怎样 |
|---|---|---|
| proposed | draft/MTP 产生,target 尚未接受 | 可被 verifier 拒绝,不得发布 |
| committed | target 已接受并进入请求逻辑序列/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 还没有内容可发送:
| tick | bounded capacity=3:sent | pending | scheduler 动作 |
|---|---|---|---|
| 0 | "" | A | 生成 A |
| 1 | "" | AB | 生成 B |
| 2 | "" | ABC | 生成 C,达到上限 |
| 3 | "" | ABC | 暂停生成 |
| 4 | A | BCD | drain A,再生成 D |
| 5 | A | BCD | 暂停生成 |
| 6 | AB | CDE | drain B,再生成 E |
| 7 | AB | CDE | 暂停生成 |
| 8 | AB | CDE | deadline;丢弃未发送并清理 |
无界路径在同一时刻已生成 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 数量近似相乘:
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 为 ,窗口 内:
是否只含 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。