NodeJS云函数中转发OpenAI流时如何拆分合并的Chunk
解决OpenAI流扩展Chunk合并问题
问题描述
我编写了extendOpenAIStream函数,目标实现:
- 接收OpenAI返回的ReadableStream
- 逐Chunk读取流数据
- 在包含
[DONE]标识的最后一个Chunk之前,插入携带自定义extensionPayload的独立Chunk
但目前遇到的问题是:插入的扩展数据与原OpenAI流的最后一个Chunk被合并成一个整体,客户端无法将它们作为独立单元处理。
原函数代码:
export async function extendOpenAIStream( openaiStream: ReadableStream<Uint8Array>, extensionPayload: JSONValue ) { const encoder = new TextEncoder() const decoder = new TextDecoder() const reader = openaiStream.getReader() const stream = new ReadableStream({ cancel() { reader.cancel() }, async start(controller) { while (true) { const { done, value } = await reader.read() const dataString = decoder.decode(value) if (done || dataString.includes('[DONE]')) { // Enque our extension const extendedValue = encoder.encode( `data: ${JSON.stringify(extensionPayload)} ` ) controller.enqueue(extendedValue) // Enque the original chunk controller.enqueue(value) // Close the stream controller.close() break } controller.enqueue(value) } }, }) return stream }
预期输出(独立Chunk):
data: {"extensionPayload": {...}} data: {"id":"...,","object":"chat.completion.chunk","created":1684486791,"model":"gpt-3.5-turbo-0301","choices":[{"delta":{},"index":0,"finish_reason":"stop"}]} data: [DONE]
实际问题:扩展数据与原最后一个Chunk的内容被合并为一个大Chunk,客户端无法区分独立单元。
问题原因
- 未处理
done状态的无效值:当流结束时done为true,value是undefined,此时调用decoder.decode(value)会得到空字符串,且controller.enqueue(value)会传入无效值。 - 未拆分包含
[DONE]的Chunk:OpenAI流的最后一个Chunk通常包含两个SSE数据块(结束的completion chunk和[DONE]块),直接原封不动转发会导致扩展数据与这两个块被客户端识别为连续流,而非独立Chunk。 - 流式解码未维护状态:未使用
TextDecoder的stream: true选项,可能导致跨Chunk的字符解码错误。
修复后的代码
export async function extendOpenAIStream( openaiStream: ReadableStream<Uint8Array>, extensionPayload: JSONValue ) { const encoder = new TextEncoder() const decoder = new TextDecoder() const reader = openaiStream.getReader() const stream = new ReadableStream({ cancel() { reader.cancel() }, async start(controller) { try { while (true) { const { done, value } = await reader.read() // 处理流结束的情况 if (done) { controller.close() break } // 流式解码,保持解码器状态 const dataString = decoder.decode(value, { stream: true }) // 检查当前Chunk是否包含[DONE]标识 if (dataString.includes('[DONE]')) { // 将原Chunk拆分为[DONE]之前的内容和[DONE]本身 const doneIndex = dataString.indexOf('[DONE]') const beforeDone = dataString.slice(0, doneIndex) const donePart = dataString.slice(doneIndex) // 1. 插入扩展数据Chunk const extendedChunk = encoder.encode( `data: ${JSON.stringify(extensionPayload)} ` ) controller.enqueue(extendedChunk) // 2. 插入[DONE]之前的原内容(如果有) if (beforeDone.trim()) { controller.enqueue(encoder.encode(beforeDone)) } // 3. 插入[DONE]部分 controller.enqueue(encoder.encode(donePart)) // 关闭流 controller.close() break } // 非最后Chunk,直接转发 controller.enqueue(value) } } catch (err) { controller.error(err) reader.cancel() } }, }) return stream }
修复说明
- 处理
done状态:当流结束时直接关闭控制器,避免处理undefined值。 - 拆分包含
[DONE]的Chunk:将原Chunk拆分为结束completion块和[DONE]块,确保扩展数据插在两者之间,三个部分均为独立Chunk。 - 流式解码:使用
decoder.decode(value, { stream: true })维护解码器内部状态,避免跨Chunk的字符解码错误。 - 错误处理:添加try-catch捕获异常,确保流能正确关闭并抛出错误。
内容的提问来源于stack exchange,提问作者flavordaaave
相关产品推荐
相关产品推荐

