You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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块内用局部变量拼接字符串,但会等待完整缓冲区解析,破坏流式实时性
解决方案

问题根源在于:

  1. 每次transform新建TextDecoder,未处理UTF-8多字节字符的碎片化(流式chunk可能是不完整的字节序列,单独解码会导致乱码或解析错误)
  2. 正则替换未考虑模式跨多个chunk的情况,直接处理单块文本会导致匹配不完整或错误
  3. 未处理流终止时的剩余缓冲区内容

以下是修正后的代码:

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.02 04:24:53