文章总结: 本文解析SGLang中TokenizerManager的实现。核心发现:利用AsyncDynamicbatchTokenizer合并并发请求进行批量编码提升效率;基于ReqState与asyncio.Event构建非阻塞异步状态管理支持流式SSE;编码后执行精准长度截断;并通过SessionParams实现多轮对话共享KVCache以节省算力。 综合评分: 86 文章分类: AI安全,代码审计,安全开发
Token 的一生:TokenizerManager —— 文本的“翻译官”
原创
zouyee zouyee
DCOS
2026年8月17日 08:00 江苏
在小说阅读器读本章
去阅读
系列定位:中级。深入
tokenizer_manager.py,理解请求的编码过程、批量优化和状态管理。
Summary
本篇深入 TokenizerManager 进程的内部实现。内容涵盖:_tokenize_one_request 的两条路径(用户传 input_ids vs 传 text);AsyncDynamicbatchTokenizer 如何合并并发请求做批量编码;InputFormat 枚举与自动格式检测;ReqState + asyncio.Event 的请求生命周期管理;流式 SSE 与非流式返回的差异;输入长度验证与自动截断;以及 SessionParams 支持多轮对话跨请求共享 KV Cache。
1. TokenizerManager 的职责
TokenizerManager 是 HTTP 服务器和调度器之间的桥梁,主要负责:
- 接收来自 HTTP 层的
GenerateReqInput - Tokenize:调用 HuggingFace tokenizer,把文本变成
array[int] - 发送
TokenizedGenerateReqInput给调度器(ZMQ PUSH) - 等待调度器/DetokenizerManager 的回包
- 流式推送或一次性返回文本给 HTTP 层
# python/sglang/srt/managers/tokenizer_manager.py
class TokenizerManager(TokenizerControlMixin, TokenizerManagerScoreMixin):
async def generate_request(self, obj, request):
obj.normalize_batch_and_arguments() # 1. 规范化
self._init_req_state(obj, request) # 2. 创建 rid_to_state 条目
tokenized_obj = await self._tokenize_one_request(obj) # 3. Tokenize
self._send_one_request(tokenized_obj) # 4. ZMQ 发送
async for response in self._wait_one_response(obj, request):
yield response # 5. 等待并流式返回
2. _tokenize_one_request 详解
async def _tokenize_one_request(self, obj):
# ① 如果已经是 input_ids(用户直接传数字),跳过 tokenize
if obj.input_ids is not None:
input_ids = obj.input_ids
else:
# ② 调用 _tokenize_texts,支持异步批处理
input_ids, token_type_ids = await self._tokenize_texts(obj.text)
# ③ 长度检查:超过最大上下文窗口则截断或报错
# ④ 封装并返回
self._validate_one_request(obj, input_ids)
return self._create_tokenized_object(
obj, input_text, input_ids, input_embeds, mm_inputs, token_type_ids
)
两条 tokenize 路径
图 1:架构与流程示意(点击图片查看原图)
异步动态批处理 tokenizer
当并发请求量很大时,每个请求单独调用 tokenizer.encode() 效率很低,因为 HuggingFace tokenizer 的快速路径(Rust 实现)在批量调用时效率更高。
AsyncDynamicbatchTokenizer 会在一个短暂的时间窗口内收集多个并发请求,合并成一个批次调用:
# 多个并发请求同时到来:
# request A: "Tell me a joke"
# request B: "What is SGLang?"
# request C: "Hello!"
#
# AsyncDynamicbatchTokenizer 把它们合并:
tokenizer(["Tell me a joke", "What is SGLang?", "Hello!"])
# 一次批量调用,比三次单独调用快 ~3x
3. 输入格式检测
用户可以传多种格式的文本,TokenizerManager 会自动判断:
class InputFormat(Enum):
SINGLE_STRING = 1 # "Tell me a joke"
BATCH_STRINGS = 2 # ["Hello", "World"]
CROSS_ENCODER_PAIRS = 3 # [["query", "document"], ...] 用于 reranker 模型
def _detect_input_format(self, texts, is_cross_encoder):
if isinstance(texts, str):
return InputFormat.SINGLE_STRING
if is_cross_encoder and isinstance(texts[0], list) and len(texts[0]) == 2:
return InputFormat.CROSS_ENCODER_PAIRS
return InputFormat.BATCH_STRINGS
Cross-encoder 格式会让 tokenizer 返回 token_type_ids,用于区分 query 和 document 的 token 归属段(segment 0 vs segment 1)。
4. ReqState:请求的生命状态
每个请求在 TokenizerManager 里都有一个对应的 ReqState:
@dataclasses.dataclass
class ReqState:
out_list: List[Dict] # 收集到的输出片段(流式下是多个)
finished: bool # 是否已完成
event: asyncio.Event # 用于 await 等待完成信号
obj: Union[GenerateReqInput, EmbeddingReqInput] # 原始请求对象
time_stats: APIServerReqTimeStats # 性能指标
# 还有 last_completion_tokens、text、text_chunks 等辅助字段
请求 ID(rid)是关键,它贯穿整个生命周期:
TokenizerManager.rid_to_state = {
"req-001": ReqState(finished=False, ...),
"req-002": ReqState(finished=False, ...),
}
# 当调度器/detokenizer 回包时:
def _handle_batch_output(self, out):
for rid, text in zip(out.rids, out.decoded_texts):
state = self.rid_to_state[rid]
state.out_list.append({"text": text})
state.event.set() # 唤醒等待该请求的协程
5. 流式与非流式的差异
非流式(stream=False)
客户端发请求 → await 所有 token 生成完毕 → 一次性返回完整文本
流式(stream=True,SSE)
客户端发请求
↓
每次 detokenizer 推来新 token 片段
↓ 立刻 yield 给客户端
↓ SSE: data: {"text": " joke"}\n\n
↓ SSE: data: {"text": " Why"}\n\n
...
↓ SSE: data: [DONE]\n\n
_wait_one_response 内部用 asyncio.Event 实现协程等待,不阻塞事件循环:
async def _wait_one_response(self, obj, request):
state = self.rid_to_state[obj.rid]
while True:
await state.event.wait() # 挂起,等待新数据到来
state.event.clear()
# 消费 state.out_list 里的数据
for item in state.out_list:
yield item
state.out_list.clear()
if state.finished:
break
6. 长度验证与截断
在发送给调度器之前,TokenizerManager 会检查 token 数量:
# 上下文窗口 = 模型支持的最大 token 数(如 128K)
# _validate_one_request 中使用 self.context_len
if input_token_num >= self.context_len:
# 根据配置决定:截断 or 报错
if self.server_args.allow_auto_truncate:
del input_ids[self.context_len:]
else:
raise ValueError(f"Input too long: {input_token_num} > {self.context_len}")
注意:截断发生在 tokenize 之后,因为只有转成 IDs 才能准确计算 token 数(”中文”两字在不同模型里可能是 1-4 个 token)。
7. Session 管理(多轮对话)
对于多轮对话,SGLang 支持 Session 模式,把多个请求关联到同一个会话:
class SessionParams:
id: Optional[str] # 会话 ID
rid: Optional[str] # 请求节点 ID(会话内的位置)
offset: Optional[int] # token 插入偏移
replace: Optional[bool] # 是否替换之前的输出
drop_previous_output: Optional[bool] # 是否丢弃上轮输出
Session 的核心优势是:同一会话的多轮请求可以共享 KV Cache,第二轮不需要重新计算第一轮的注意力权重,大幅节省计算量。这在 RAG、agent 循环等场景下效果显著。
小结
| 环节 | 关键代码 | 作用 |
| — | — | — |
| 文本编码 | _tokenize_texts() | 文本 → array[int] |
| 批量优化 | AsyncDynamicbatchTokenizer | 并发请求合并批量编码 |
| 状态管理 | rid_to_state | 追踪每个请求的生命周期 |
| 等待机制 | asyncio.Event | 协程挂起,不阻塞事件循环 |
| 多轮对话 | SessionParams | 跨请求共享 KV Cache |
下一篇:进入调度器,看 Token 预算是怎么分配的,以及 Chunked Prefill 如何让长提示不阻塞短请求。
免责声明:
本文所载程序、技术方法仅面向合法合规的安全研究与教学场景,旨在提升网络安全防护能力,具有明确的技术研究属性。
任何单位或个人未经授权,将本文内容用于攻击、破坏等非法用途的,由此引发的全部法律责任、民事赔偿及连带责任,均由行为人独立承担,本站不承担任何连带责任。
本站内容均为技术交流与知识分享目的发布,若存在版权侵权或其他异议,请通过邮件联系处理,具体联系方式可点击页面上方的联系我。
本文转载自:DCOS zouyee zouyee《Token 的一生:TokenizerManager —— 文本的“翻译官”》
版权声明
本站仅做备份收录,仅供研究与教学参考之用。
读者将信息用于其他用途的,全部法律及连带责任由读者自行承担,本站不承担任何责任。










评论