Prhub

#50868 [Rust][Benchmark] Preserve UTF-8 across benchmark stream chunks

原始 PR 作者 reidliu41 合并时间 2026-08-04 14:15 文件变更 1 提交数 1 评论 1 代码增减 +30 / -10

执行摘要

修复 Rust benchmark SSE 跨 chunk UTF-8 乱码

PR body 明确说明:原 SSE handler 对每个网络 chunk 独立执行 String::from_utf8_lossy,当 TCP 或 HTTP chunking 拆分多字节 UTF-8 字符(如 中 的 E4 B8 AD 在 E4 后被切)时,每个不完整字节序列会在后续 chunk 到达前被替换为 U+FFFD,流式文本被破坏。需要改为先以原始字节缓存、到完整 SSE 消息再解码,保证非 ASCII 文本的 benchmark 输出完整,避免损坏的响应被保存或作为多轮对话历史复用;同时明确要求对真正非法 UTF-8 的 lossy 行为保持不变。

值得快速精读:PR 虽小,但体现了一个可复用的流式协议处理模式——先累积原始字节、再按完整消息边界解码,并用 Cow 临时视图做投机解析而不破坏字节缓冲。对 Rust 后端开发者和 benchmark 维护者都有参考价值;建议关注 SSE buffer 在字节层面切分与 UTF-8 解码时机的权衡,以及 windows(2) 查找与 String::find 在语义上的等价性。

讨论亮点

本 PR 没有实质技术讨论。esmeetu 直接回复 LGTM! 并批准;claude[bot] 说明这是 fork 提交,自动 review 被禁用,维护者可用 @claude review 触发一次性 review。技术决策主要依靠 PR body 中的自述与测试用例呈现,未出现设计权衡争议。

实现拆解

  1. 变更入口:rust/src/bench/src/backends/streaming.rs 中的 StreamedResponseHandler,这是 Rust benchmark 工具处理 SSE 流式响应(streaming backend)的唯一核心结构。
  2. 缓冲从 String 改为 Vec:struct 字段 buffer 由 String 改为 Vec,new() 中改为 Vec::with_capacity(4096)。这样可以在不损失性能的前提下累积任意字节序列,不再在增量阶段强制解码。
  3. add_chunk 先累积后解码:原先每次调用先 String::from_utf8_lossy(chunk_bytes) 再 push_str;现在改为 self.buffer.extend_from_slice(chunk_bytes),将原始字节原样加入缓冲区。切分 SSE 消息时改用 `windows(2).position(|window| window == b"

")在字节层面定位双换行分隔符,切出的完整消息段才调用String::from_utf8_lossy(...).trim().to_string()解码;self.buffer.drain(..pos + 2)` 的消费逻辑保持不变。

  1. 投机 JSON 解析改用临时视图:尾部无 `

的 speculative 解析路径(对应 Python 端 StreamedResponseHandler 的 speculative json.loads(),对 TTFT/ITL 指标准确性很重要)现在通过let buffer = String::from_utf8_lossy(&self.buffer)创建临时 Cow 视图,data_start 查找、content 解析、messages.push(buffer.trim().to_string())都基于该视图;解析成功后self.buffer.clear()` 清空原始字节。不完整的多字节字符字节仍留在 Vec 中参与下一 chunk 的累积。

  1. 测试配套:在 mod tests 内新增 test_split_multibyte_utf8(含 CJK 与 emoji 的消息逐字节投喂,断言只产出一条完整消息)和 test_invalid_utf8_remains_lossy(0xFF 字节仍输出 U+FFFD,保证 lossy 行为不变)。测试随源码文件一同提交,未单独拆出测试文件。
文件 模块 状态 重要度
rust/src/bench/src/backends/streaming.rs 流式解析 modified 6.99

关键符号

StreamedResponseHandler::new StreamedResponseHandler::add_chunk test_split_multibyte_utf8 test_invalid_utf8_remains_lossy

关键源码片段

rust/src/bench/src/backends/streaming.rs core-logic

唯一变更文件,承载全部核心逻辑:StreamedResponseHandler 的 buffer 从 String 改为 Vec<u8>,add_chunk 不再按 chunk 独立解码,改为累积原始字节并在完整 SSE 消息边界解码;投机 JSON 解析改用临时 lossy 视图。同文件内新增两个针对性测试。

// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright contributors to the vLLM project/// SSE 流式响应处理器。
///
/// 以原始字节累积网络 chunk,只在完整 SSE 消息边界才做 UTF-8 解码,
/// 避免 TCP 或 HTTP 分片把多字节字符切成两半时,前一半立刻被
/// `String::from_utf8_lossy` 替换为 U+FFFD,导致流式文本乱码。
pub struct StreamedResponseHandler {
    // buffer 存原始字节而不是 String:跨 chunk 的多字节字符保持完整
    buffer: Vec<u8>,
    /// 可复用消息缓冲,避免每次 `add_chunk` 都新分配 Vec
    messages: Vec<String>,
}impl StreamedResponseHandler {
    pub fn new() -> Self {
        Self {
            buffer: Vec::with_capacity(4096),
            messages: Vec::with_capacity(4),
        }
    }    /// 加入一个字节 chunk,返回完整 SSE 消息。
    ///
    /// 返回值借用 `self`,只在下次 `add_chunk` 前有效。
    pub fn add_chunk(&mut self, chunk_bytes: &[u8]) -> &[String] {
        self.messages.clear();        // 1. 先累积原始字节,不做任何解码,保证不完整多字节字符跨 chunk 保留
        self.buffer.extend_from_slice(chunk_bytes);        // 2. 按 SSE 消息分隔符(双换行)切分,切出的完整消息才解码;
        // 这里在字节层面用 windows(2) 定位 b"\n\n",等价于原 String 里 find
        while let Some(pos) = self.buffer.windows(2).position(|window| window == b"\n\n") {
            let message = String::from_utf8_lossy(&self.buffer[..pos]).trim().to_string();
            // 移除已消费字节,剩余数据前移
            self.buffer.drain(..pos + 2);
            if !message.is_empty() {
                self.messages.push(message);
            }
        }        // 3. 尾部没有 \n\n 的“投机 JSON 解析”路径,对齐 Python 端
        // StreamedResponseHandler 里的 speculative json.loads(),保证
        // TTFT/ITL 指标准确性:data 消息与其结束分隔符分属两个 TCP 段时,
        // 可以在第一个段到达时上报,而不是等第二个段。
        //
        // 这里用 `String::from_utf8_lossy` 构建临时 Cow 视图用于查找和解析;
        // 原始字节仍留在 `self.buffer`,不完整的多字节字符可等下一个 chunk 补齐。
        let buffer = String::from_utf8_lossy(&self.buffer);
        let data_start = if buffer.starts_with("data: ") {
            Some(0)
        } else {
            // 兼容多字段 SSE 事件(如 Dynamo 前端发的 "event: ...\ndata: ...")
            buffer.find("\ndata: ").map(|p| p + 1)
        };
        if let Some(offset) = data_start {
            let content = buffer[offset + 6..].trim();
            if content == "[DONE]"
                || (!content.is_empty()
                    && serde_json::from_str::<&serde_json::value::RawValue>(content).is_ok())
            {
                self.messages.push(buffer.trim().to_string());
                // 已消费整条消息,清空累积的原始字节
                self.buffer.clear();
            }
        }        &self.messages
    }
}
#[cfg(test)]
mod tests {
    use super::*;    #[test]
    fn test_split_multibyte_utf8() {
        // 逐字节投喂含 CJK 与 emoji 的 SSE 消息,端到端结果必须与一次性投喂
        // 完全一致,证明跨 chunk 的 UTF-8 序列没有被破坏。
        // 中 是 3 字节(E4 B8 AD),😀 是 4 字节,chunks(1) 会任意切割。
        let mut handler = StreamedResponseHandler::new();
        let message = "data: {\"test\":\"中😀\"}\n\n";
        let mut messages = Vec::new();        for chunk in message.as_bytes().chunks(1) {
            messages.extend_from_slice(handler.add_chunk(chunk));
        }        assert_eq!(messages, &[message.trim().to_string()]);
    }    #[test]
    fn test_invalid_utf8_remains_lossy() {
        // 真正非法的 UTF-8 字节仍按原行为替换为 U+FFFD,不因本修复改变
        let mut handler = StreamedResponseHandler::new();
        let msgs = handler.add_chunk(b"data: {\"test\":\"\xFF\"}\n\n");
        assert_eq!(msgs, &["data: {\"test\":\"\"}".to_string()]);
    }
}

评论区精华

审核流程与自动 review 状态 other

claude[bot] 注明本 PR 来自 fork、自动 review 被禁用,维护者可评论 @claude review 触发一次性 review;esmeetu 直接回复 "LGTM!" 并批准合并。

结论:无技术异议,维护者认可后直接合并,未产生设计或正确性交锋。 · 已解决

风险与影响

风险面较窄:

1) 变更仅涉及 benchmark 工具的 SSE 流式后端,不触及 vLLM 推理主链路、KV cache 或调度等核心路径,回归面小。
2) 行为一致性风险:投机 JSON 解析路径改为基于临时 lossy 视图,若视图内容与最终完整解码结果不一致(例如多字节字符跨消息边界),可能导致提前触发消息上抛的时序变化,但该路径逻辑与原实现保持同构,仅解码时机后移。
3) 一个理论上未覆盖的边界:若多字节字符字节序列与 `

分隔符交错到达,windows(2)可能在多字节字符中间切分,该半截字符仍会在该消息解码时变成 U+FFFD;但 SSE 帧内 JSON 字符串不允许裸换行,

` 只会出现在结构化分隔位置,实际触发概率极低,这一点在材料中未被讨论与测试覆盖。

4) Vec 累积与原有 String 累积的内存特性基本一致,无可见性能回退。

影响范围:仅影响使用 rust bench(vllm-bench)进行流式 benchmark 的用户与 CI 场景,修复的是非 ASCII 输出被损坏的可见缺陷,避免错误响应被保存或复用为多轮对话上下文。对 vLLM 推理运行时无影响。对团队而言,这是 Rust 工具链正确性维护的一部分,代码量小、风险低,但保障了 benchmark 数据的可信度。

流式解析路径变更 影响范围限于 benchmark 工具 投机解析路径存在未覆盖边界

关联 Issue

未识别关联 Issue

当前没有检测到明确关联的 Issue 链接,后续同步到相关引用后会出现在这里。

完整报告

参与讨论