解决流式LLM响应「必须等完整JSON才能解析」的痛点,构建部分JSON读取器在每个chunk到达时返回当前最完整的值,并指出两种常见实现导致数据静默损坏的坑。
流式 LLM JSON 解析:在最后一个 Token 到来之前
一个能可靠返回有效 JSON 的模型,在打字的过程中并不会返回有效 JSON。这是两个完全不同的保证,而大多数流式客户端会悄悄地把它们当作同一件事。约束解码和 schema 强制输出模式承诺最终响应可以解析。它们不会承诺第十七个 chunk 会是什么——它可能结束在一个 key 的中间、一个转义序列的内部,或者紧跟在一个还没有后续元素的逗号之后。
所以常见的妥协方案是:把 token 流出来供一个进度旋转动画使用,然后无论如何都全部缓冲起来,最后一次性解析。用户看着字符一个个出现,而应用程序自己却拒绝读取。模型已经决定的所有内容——标题、前三个列表项、在前五十个 token 中就到达的分类标签——都只能放在字符串缓冲区里无法使用,直到闭合花括号落地。
这个差距是可以弥合的。下面要构建的是一个部分 JSON 读取器,它把每一个 chunk 转换为流在当前条件下能合理确认的完整值,然后指出这个思路的明显版本会在两个地方悄悄破坏数据。
流式补全以文本增量(delta)的形式到达,与 JSON 结构毫无关系。Token 边界遵循分词器的词表,而不是语法。一个单独的 delta 可能携带 {"ti,或者 "tle": "Sh,或者一个孤立的反斜杠——只有等下一个 delta 提供了 n 之后它才有意义。
客户端循环很简单:
import json
from openai import OpenAI
client = OpenAI()
def stream_text(prompt: str):
response = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
response_format={"type": "json_object"},
stream=True,
)
for chunk in response:
delta = chunk.choices[0].delta.content
if delta:
yield delta
这个生成器的每个元素都是一份还不成文档的文档片段。问题是,消费者对每个 chunk 除了追加之外还能做什么。
标准库解析器是全有或全无的设计。给它一个前缀它就会抛出异常,因为 JSON 文档的前缀本身不是 JSON 文档。
buffer = ""
for delta in stream_text(prompt):
buffer += delta
try:
value = json.loads(buffer) # 除了最后一个 chunk 外每个都会 raise
except json.JSONDecodeError:
continue
json.JSONDecoder.raw_decode 看起来像是一条逃生通道,而且它确实有用,但解决的是不同的问题:它从字符串开头解码一个完整值,并报告停在了哪里。这解决了流中分离对象的拼接文档问题。但当你想取的单个对象被截断时它就帮不上忙了,因为开头仍然没有一个完整值可以解码。
ijson 包更接近需求。它是一个事件驱动的解析器,在字节到达时产生 start_map、map_key、string 和 end_map 事件,这正是你想要的流式形态。它的约束是它为非常大但最终会完整的文档构建的;一个被截断的 feed 会抛出 IncompleteJSONError,你只能得到已经关闭的值的事件。对于一个想在字符串增长时展示它的渲染 UI 来说,每个完成值一个事件还是太粗糙了。
写一个真正的增量解析器是彻底的答案,但它的工作量超出了问题本身的价值。一个分词器、一个基于语法状态机、以及一个可以在中途被查询的值构建器,需要几百行代码来编写,而且要信任它还需要更多。剩下的方案更便宜,而且复用一个已经正确的解析器:把前缀修复成有效文档,解析这个修复后的结果,然后把修复后的内容丢弃。缓冲区本身永远不会被修改,所以修复逻辑中的一个错误只付出一个糟糕快照的代价,而不是污染整个流。
核心思想很小。追踪哪些容器是开放的,当请求快照时,追加能够完成它们的闭合符。
def naive_snapshot(buffer: str):
stack = []
for ch in buffer:
if ch == "{":
stack.append("}")
elif ch == "[":
stack.append("]")
elif ch in "}]" and stack:
stack.pop()
try:
return json.loads(buffer + "".join(reversed(stack)))
except json.JSONDecodeError:
return None
对于一个表现良好的前缀,这立即就能工作。{"title": "Ship it", "tags": ["python" 变成 {"title": "Ship it", "tags": ["python"]} 并解析成一个模板可以立即渲染的字典,比响应结束早了几百个 token。
它也会不断失败,而且失败比成功更有意思。{"tags": ["python", 闭合成一个尾部逗号。{"score": 闭合成一个没有值的 key。{"score": 1. 闭合成一个以小数点结尾的数字字面量。每一个都是解码错误,所以快照返回 None,UI 就会停滞,直到流恰好落在一个幸运的边界上。
危险的失败与那些嘈杂的失败不同。嘈杂的失败会自我宣告:尾部逗号抛出异常,快照是 None,下一个 chunk 通常会修复它。字符串值内部的括号不是括号,不会抛出任何异常。
考虑前缀 {"note": "use {curly} braces。上面的扫描把 prose 中的 { 也计数了,向栈中推入第二个 },永远不会看到匹配的闭合,产生 {"note": "use {curly} braces"}}。这不是一个解码错误。取决于流在哪里停止,这样的修复可能产生一个解析成错误形状的文档——一种从模型 prose 中臆造出来的结构。
这种输入也不罕见。任何在谈论代码、文件路径或模板语法的模型都会在字符串值内部不断发出括号和方括号,而关于 JSON 的 JSON 响应是朴素扫描器最糟糕的情况。一个解析成错误形状的快照比解析失败的快照更糟糕,因为消费者没有任何信号表明出了问题。
任何不知道自己在字符串内部的扫描器都是在猜测。同样的道理适用于转义序列:以单个 \ 结尾的缓冲区是一个不完整的转义序列,用 " 关闭字符串会把尾部反斜杠变成一个转义的引号,这会吞掉闭合符,并把腐败向外推一层。
正确的处理需要在 chunk 之间携带三个状态:容器栈、一个 in-string 标志和一个 escape 标志。
每个字符在到达时被分类一次,永远不会被重新扫描。每一帧也记录一个安全截断偏移量——容器只包含完整元素的位置——这样当修复失败时可以退回到最后一个已知的好边界,而不是什么都返回。
import json
from dataclasses import dataclass
@dataclass
class _Frame:
close: str
cut: int # offset where this container held only finished elements
class PartialJSONStream:
def __init__(self) -> None:
self._buf: list[str] = []
self._len = 0
self._stack: list[_Frame] = []
self._in_string = False
self._escaped = False
def feed(self, chunk: str) -> None:
for ch in chunk:
self._step(ch)
self._buf.append(ch)
self._len += 1
def _step(self, ch: str) -> None:
i = self._len
if self._in_string:
if self._escaped:
self._escaped = False
elif ch == "\\":
self._escaped = True
elif ch == '"':
self._in_string = False
return
if ch == '"':
self._in_string = True
elif ch == "{":
self._stack.append(_Frame("}", i + 1))
elif ch == "[":
self._stack.append(_Frame("]", i + 1))
elif ch in "}]":
if self._stack:
self._stack.pop()
elif ch == "," and self._stack:
self._stack[-1].cut = i
快照方法从最完整到最保守生成修复候选,并返回第一个能解析的:
def _closers(self, upto: int | None = None) -> str:
frames = self._stack if upto is None else self._stack[:upto]
return "".join(frame.close for frame in reversed(frames))
def _candidates(self, text: str):
head = text[:-1] if self._escaped else text
if self._in_string:
head += '"'
head = head.rstrip()
if head.endswith(","):
head = head[:-1].rstrip()
elif head.endswith(":"):
head += " null"
yield head + self._closers()
for depth in range(len(self._stack) - 1, -1, -1):
cut = self._stack[depth].cut
yield text[:cut].rstrip().rstrip(",") + self._closers(depth + 1)
def snapshot(self):
text = "".join(self._buf)
for candidate in self._candidates(text):
try:
return json.loads(candidate)
except json.JSONDecodeError:
continue
return None
回退路径使得部分数字和半输入的字面量变得无害。以 1. 或 tru 结尾的缓冲区产生的第一个候选会失败,然后一个截断的候选会完全丢弃未完成的元素并解析。调用者看到的是没有那个 key 的对象,而不是什么都看不到。
必须对消费的代码明确说明这个设计的一个特性:快照中的字符串值可能是真实值的前缀。在增长时渲染一个部分描述是预期的用途。将部分值与枚举比较、将其作为 URL 使用、或传递给执行某个动作的东西,都是等待慢 token 的 bug。
每个 chunk 都轮询整个快照,对于只关心字段何时完成的消费者来说是浪费的。JSON 对象按顺序流式传输,这给出了一个简单可靠的完成规则:一旦出现第二个 key,第一个 key 的值就不可能再改变了。快照中除了最后一个 key 之外的每个 key 都是已确定的。
class FieldEmitter:
def __init__(self) -> None:
self.stream = PartialJSONStream()
self._emitted: set[str] = set()
def feed(self, chunk: str) -> list[tuple[str, object]]:
self.stream.feed(chunk)
snap = self.stream.snapshot()
if not isinstance(snap, dict):
return []
events = []
for key in list(snap)[:-1]:
if key not in self._emitted:
self._emitted.add(key)
events.append((key, snap[key]))
return events
def finish(self) -> list[tuple[str, object]]:
snap = self.stream.snapshot()
if not isinstance(snap, dict):
return []
return [(k, v) for k, v in snap.items() if k not in self._emitted]
下游处理器现在可以在模型还在写解释时就开始进行分类工作,这才是流式传输结构化响应而不是一团 prose 的意义所在。路由决策可以被分发、数据库行可以被预留、UI 区域可以渲染为最终状态而不是骨架状态。
这个排序规则有一个附加条件:它适用于模型以固定 key 顺序写入的对象,这就是 schema 约束解码产生的结果。它不适用于对象数组中后面的元素修改了前面元素的情况,所以把列表的最后一个元素当作与对象的最后一个 key 完全相同的方式视为暂定的。
快照是一个草稿,所以针对严格的输出模型验证它是结构上的错误——必填字段按预期是缺失的。为快照构建一个宽松的模型镜像,把严格的模型留给最终值。
from pydantic import BaseModel, create_model
class Review(BaseModel):
verdict: str
reasons: list[str]
score: int
def draft_of(model: type[BaseModel]) -> type[BaseModel]:
fields = {
name: (info.annotation | None, None)
for name, info in model.model_fields.items()
}
return create_model(f"Draft{model.__name__}", **fields)
DraftReview = draft_of(Review)
有两个失败模式需要围绕它进行明确处理。第一个是模型停止产生结构而开始产生道歉或带围栏的代码块;快照变成 None 并一直保持 None,所以追踪连续无法解析的 chunk 并放弃流,而不是等待一个不会到来的闭合。第二个是提前退出:当快照已经满足严格模型且剩余字段是可选的时,取消请求会停止支付没有人会读的 token 的费用。
还要对缓冲区设置上限。部分解析器会乐意地从陷入重复循环的模型中累积兆字节,缓冲区的大小上限是防止这种情况的最便宜的保护。
值得写的测试是以最差的粒度输入的测试,因为 chunk 大小为 1 会锻炼真实分词器可能产生的每一个边界。
import json
import pytest
from partial_json import PartialJSONStream
PAYLOADS = [
{"verdict": "ship", "reasons": ["tests pass", "small diff"], "score": 8},
{"note": 'use {curly} braces and a "quote"', "path": "C:\\tmp\\out.json"},
{"nested": {"a": [1, 2, {"b": None}], "c": True}, "trailing": 1.5},
]
@pytest.mark.parametrize("payload", PAYLOADS)
def test_every_prefix_parses_or_declines(payload):
text = json.dumps(payload)
stream = PartialJSONStream()
for ch in text:
stream.feed(ch)
snap = stream.snapshot()
assert snap is None or isinstance(snap, (dict, list))
assert stream.snapshot() == payload
@pytest.mark.parametrize("payload", PAYLOADS)
def test_settled_keys_never_change(payload):
text = json.dumps(payload)
stream = PartialJSONStream()
seen: dict[str, object] = {}
for ch in text:
stream.feed(ch)
snap = stream.snapshot()
if isinstance(snap, dict):
for key in list(snap)[:-1]:
if key in seen:
assert seen[key] == snap[key]
else:
seen[key] = snap[key]
第二个测试才是重要的。它编码了 emitter 向消费者做出的承诺——一个作为已确定值发出的值永远不会被修改——而且它是能够捕获扫描器 bug 的测试,这种 bug 只在括号出现在字符串内部时才会显现,因为被破坏的快照会改变一个已经报告过的 key。
基于属性的测试可以廉价地扩展这一点。用 Hypothesis 生成任意嵌套结构,序列化它们,喂入每个前缀,然后断言同样的两个不变量。任何转义处理错误都会作为缩小后的反例浮现,而不是作为关于字段被破坏的支持工单。
这不是一个更快的模型。它移除了一个从未必要的等待——一个值被确定下来的时刻和闭合花括号让应用程序承认它已经知道的时刻之间的间隔。解析器大约一百行代码,它携带的状态是三个变量和一个栈,正确性论证可以放进两个测试里。这是在答案还在被写入时就展示给用户的合理代价。
最初发表于 Dispatch。