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

如何在每次pull时刷新ReadableStream内部队列?

问题分析与解决方案

你的核心问题出在客户端流读取逻辑不完整,以及流数据处理的编码/拆分逻辑错误,导致看起来只有迭代器结束时才会收到所有数据。以下是具体修正方案:

1. 服务端代码微调(确保分块传输)

在返回Response时明确指定分块传输头,避免浏览器缓存数据直到流结束:

function iteratorToStream(iterator: AsyncGenerator<Uint8Array, void, undefined>) {
  const readStream: ReadableStream<Uint8Array> = new ReadableStream(
    {
      async pull(controller) {
        const { value, done } = await iterator.next()
        if (done) {
          controller.close()
        } else {
          controller.enqueue(value)
        }
      },
    },
    {
      highWaterMark: 1,
    }
  )
  return readStream;
}

const encoder = new TextEncoder()
// 补充sleep函数实现
function sleep(ms: number) {
  return new Promise(resolve => setTimeout(resolve, ms));
}

async function* makeIterator() {
  yield encoder.encode('HelloWorldHelloWorldHelloWorldHelloWorldHelloWorld\n');
  yield encoder.encode('HelloWorldHelloWorldHelloWorldHelloWorldHelloWorld\n');
  await sleep(1000);
  yield encoder.encode('延迟1秒后的消息\n');
}

export async function POST() {
  const iterator = makeIterator()
  const stream = iteratorToStream(iterator)

  return new Response(stream, {
    headers: {
      'Content-Type': 'text/plain; charset=utf-8',
      'Transfer-Encoding': 'chunked' // 明确启用分块传输
    }
  })
}

2. 客户端代码修正(实现逐消息读取)

客户端存在三个关键错误:只调用一次reader.read()、未正确解码Uint8Array、未按分隔符拆分消息。修正后的代码如下:

function logStream(splitOn: string) {
  let buffer = '';
  const decoder = new TextDecoder(); // 用于解码Uint8Array为字符串

  return new TransformStream({
    transform(chunk, controller) {
      // 把二进制chunk解码为字符串,stream: true表示支持流式解码
      buffer += decoder.decode(chunk, { stream: true });
      // 按指定分隔符拆分消息
      const parts = buffer.split(splitOn);
      // 将完整的消息入队,保留最后一段未完成的内容在buffer中
      for (let i = 0; i < parts.length - 1; i++) {
        if (parts[i].trim()) { // 跳过空行
          console.log(`收到完整消息: ${parts[i]}`);
          controller.enqueue(parts[i]);
        }
      }
      buffer = parts[parts.length - 1];
    },
    flush(controller) {
      // 处理流结束时剩余的未完成消息
      if (buffer.trim()) {
        console.log(`收到最后一条消息: ${buffer}`);
        controller.enqueue(buffer);
      }
    }
  });
}

async function consumeStream() {
  const response = await fetch('/route');
  if (!response.body) throw new Error('响应体为空');

  const reader = response.body
    .pipeThrough(logStream('\n')) // 按换行符拆分消息
    .getReader();

  // 循环读取流,直到所有数据接收完成
  while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    // 这里可以添加每条消息的业务处理逻辑
  }
}

// 启动流消费
consumeStream();

关键修正点说明

  • 分块传输头:服务端添加Transfer-Encoding: chunked,告诉浏览器不要缓存数据,收到一块就返回一块。
  • 循环读取流:客户端必须循环调用reader.read(),因为单次read()只能获取一个chunk,循环才能持续接收后续的流数据。
  • 正确解码与拆分:使用TextDecoder将Uint8Array转为字符串,再按指定分隔符(这里是换行符)拆分,确保每条yield的内容被识别为单独的消息。

内容的提问来源于stack exchange,提问作者Issac Spiegel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 18:03:19