如何在Next.js App Router API路由中中途处理LLM流式响应?
问题:Next.js自定义TransformStream处理流式响应时出现异常
直接使用适配器流式传输数据一切正常,但通过自定义TransformStream拦截块缓冲区、解码文本、运行正则表达式并重新编码时,流会出现锁定、随机丢失文本块,或客户端抛出以下错误:
TypeError: ReadableStream pipeline crashed or terminated unexpectedly mid-flight.
代码配置(app/api/chat/route.ts)
import { NextResponse } from 'next/server'; export async function POST(req: Request) { const { messages } = await req.json(); // 假设的AI提供商流式响应 const aiStream = await callAIProviderStream(messages); // 用于中途修改块的自定义转换流 const transformStream = new TransformStream({ transform(chunk, controller) { const text = new TextDecoder().decode(chunk); // 尝试在推送给客户端前修改文本字符串 const cleanedText = text.replace(/\[source:\s*\d+\]/g, ''); controller.enqueue(new TextEncoder().encode(cleanedText)); } }); const mutatedStream = aiStream.pipeThrough(transformStream); return new NextResponse(mutatedStream, { headers: { 'Content-Type': 'text/event-stream' }, }); }
已尝试的方法
- 验证未使用pipeThrough时原始aiStream运行正常
- 尝试在transform块内用局部变量拼接字符串,但会等待完整缓冲区解析,破坏流式实时性
解决方案
问题根源在于:
- 每次transform新建
TextDecoder,未处理UTF-8多字节字符的碎片化(流式chunk可能是不完整的字节序列,单独解码会导致乱码或解析错误) - 正则替换未考虑模式跨多个chunk的情况,直接处理单块文本会导致匹配不完整或错误
- 未处理流终止时的剩余缓冲区内容
以下是修正后的代码:
import { NextResponse } from 'next/server'; export async function POST(req: Request) { const { messages } = await req.json(); const aiStream = await callAIProviderStream(messages); // 复用TextDecoder,开启stream模式处理碎片化字节 const decoder = new TextDecoder('utf-8', { stream: true }); // 维护缓冲区处理跨chunk的正则匹配 let buffer = ''; // 要替换的正则表达式(全局匹配) const sourceRegex = /\[source:\s*\d+\]/g; const transformStream = new TransformStream({ transform(chunk, controller) { // 解码当前chunk,结合之前的不完整序列 buffer += decoder.decode(chunk); // 循环匹配所有完整的目标模式 let match; let lastIndex = 0; while ((match = sourceRegex.exec(buffer)) !== null) { // 发送匹配前的文本 if (match.index > lastIndex) { controller.enqueue(new TextEncoder().encode(buffer.slice(lastIndex, match.index))); } // 跳过匹配的内容(即替换为空) lastIndex = sourceRegex.lastIndex; } // 保留未匹配完成的尾部内容(避免跨chunk的模式被截断) buffer = buffer.slice(lastIndex); }, // 流终止时处理剩余缓冲区 flush(controller) { // 解码剩余的字节(如果有的话) buffer += decoder.decode(); // 发送剩余的文本 if (buffer) { controller.enqueue(new TextEncoder().encode(buffer.replace(sourceRegex, ''))); } } }); const mutatedStream = aiStream.pipeThrough(transformStream); return new NextResponse(mutatedStream, { headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', 'Connection': 'keep-alive' }, }); }
关键修正点说明
- 复用TextDecoder并开启stream模式:
new TextDecoder('utf-8', { stream: true })会自动保留不完整的UTF-8字节序列,直到下一次解码时拼接完整,避免多字节字符碎片化导致的乱码或解析错误。 - 维护文本缓冲区:通过
buffer变量积累文本,确保跨chunk的正则模式能被完整匹配,避免遗漏或错误替换。 - flush方法处理剩余内容:在流结束时,处理缓冲区中剩余的文本,确保所有内容都被处理并发送到客户端。
- 优化响应头:添加
Cache-Control和Connection头,确保客户端正确处理流式响应,避免缓存或连接中断问题。
内容的提问来源于stack exchange,提问作者Ashikur Rahman
相关产品推荐
相关产品推荐

