1. 从轮询到流式:为什么我们需要“流”?
如果你在十年前做前端,处理服务器推送数据,大概率会跟“轮询”(Polling)和“长轮询”(Long Polling)打交道。那时候,为了在网页上实现一个实时更新的股票价格或者聊天消息,前端得隔几秒就发个请求去问服务器:“有数据吗?” 服务器说:“没有。” 过几秒再问。这种方式笨重、低效,浪费服务器和网络资源,用户体验也差,总有延迟。
后来,我们有了Server-Sent Events,也就是常说的SSE。这玩意儿本质上是一个长连接,服务器可以主动、持续地向客户端推送数据流。前端用EventSourceAPI 来接收,感觉就像打开了一个水龙头,数据源源不断地流过来。这比轮询优雅太多了,一度成为实时数据展示(比如监控仪表盘、新闻推送)的首选方案。
但EventSource也有它的局限。它是一个“只读”的流,协议相对固定(只能是 text/event-stream 格式),而且一旦连接建立,你就只能被动接收,很难在中间对数据进行复杂的处理,比如解密、转换格式、或者与来自其他源的数据流进行合并。它更像一个设计好的管道,你只能从另一端接水喝。
于是,更强大、更底层的Web Streams API出现了。它提供了一整套处理流数据的原生工具,其中ReadableStream和TransformStream是两个核心角色。ReadableStream代表了数据的源头,你可以用各种方式去“读取”它;TransformStream则像一个安装在管道中间的“处理器”或“过滤器”,可以对流过的数据进行变换。
所以,标题里说的“三次进化”,并不是说后者完全取代前者,而是代表了前端处理流式数据能力的三次范式升级:
- EventSource: 提供了开箱即用的、基于 HTTP/1.1 的服务器推送能力,简单但受限。
- ReadableStream: 提供了底层的、通用的流读取能力,解放了数据源,让我们可以处理任何类型的流(包括 fetch 响应、本地文件、甚至自定义生成的数据)。
- TransformStream: 在
ReadableStream的基础上,增加了强大的中间处理能力,让数据流在传输过程中可以被灵活地转换、加工,实现了更复杂的流式编程模式。
今天,我们就来深入聊聊这三者,不只是看 API 怎么用,更要弄明白它们各自解决了什么问题,在什么场景下该选谁,以及如何组合它们来构建更健壮的实时应用。我会结合很多实际踩坑的经验,告诉你哪些地方容易出问题,以及怎么避开。
2. EventSource:简单直接的服务器推送
我们先从最熟悉的EventSource开始。它的 API 简单到令人发指,这也是它早期快速普及的原因。
2.1 基本用法与核心机制
假设你的服务器提供了一个 SSE 端点,比如https://api.example.com/events。前端只需要几行代码:
const eventSource = new EventSource('https://api.example.com/events'); // 监听指定类型的事件 eventSource.addEventListener('stock-update', function(event) { const data = JSON.parse(event.data); console.log('股票更新:', data); }); // 监听默认的 'message' 事件(如果服务器发送的事件没有指定 type,或者 type 为 'message') eventSource.onmessage = function(event) { console.log('收到消息:', event.data); }; // 监听连接打开事件 eventSource.onopen = function() { console.log('连接已建立'); }; // 监听错误事件 eventSource.onerror = function(error) { console.error('连接错误:', error); // 注意:EventSource 在出错时会自动尝试重连 };服务器端的响应必须遵循 SSE 格式:Content-Type: text/event-stream,并且数据以特定格式发送。一个典型的服务器响应体看起来像这样:
data: {"symbol":"AAPL","price":175.32} event: stock-update data: {"symbol":"GOOGL","price":135.67} : 这是一条注释行,客户端会忽略 data: 这是一条多行消息的第一行 data: 这是第二行关键点解析:
data:: 表示一个数据行。多个data:行会连接成一个完整的event.data字段,用换行符分隔。event:: 指定事件类型。对应前端addEventListener监听的事件名。如果省略,则默认为'message'。id:: 设置事件 ID,用于断线重连时,客户端可以通过Last-Event-ID头告诉服务器“我从哪个 ID 之后的数据开始要”。retry:: 指定重连时间(毫秒)。这在网络不稳定时非常有用。
2.2 优势与适用场景
EventSource的优势非常明显:
- 极其简单: 客户端 API 直观,服务器端实现也相对容易(几乎所有后端语言都有库支持)。
- 自动重连: 内置了连接失败后的自动重连机制,并且通过
Last-Event-ID支持断点续传,这对稳定性要求高的应用是福音。 - 基于 HTTP/1.1: 不需要 WebSocket 那样的协议升级,兼容性极好,穿透大多数防火墙和代理没问题。
它最适合那些“服务器主动发,客户端被动收”的场景:
- 实时通知: 新闻推送、系统公告、版本更新提示。
- 监控仪表盘: 服务器状态监控(CPU、内存)、业务指标(在线用户数、订单量)的实时图表更新。
- 简单的进度报告: 长任务执行进度的反馈。
2.3 局限性:为什么我们需要进化?
尽管好用,但EventSource的局限性在复杂应用面前逐渐暴露:
- 只读文本流: 它只能处理
text/event-stream格式的文本数据。如果你想传输二进制数据(如图片、音频片段),或者使用其他格式(如 ndjson),就需要在客户端额外进行 Base64 编解码或手动解析,增加了复杂性和开销。 - 单向通信: 客户端无法通过同一个连接向服务器发送数据。虽然你可以用另一个 HTTP 请求来实现,但这破坏了“流”的上下文一致性,并且增加了复杂度。
- 有限的流控制: 你无法控制流的流速。如果客户端处理不过来,数据会在缓冲区堆积,可能导致内存溢出。
EventSource没有内置的反压(Backpressure)机制。 - 难以组合与转换: 你无法轻松地将一个
EventSource流出来的数据,通过一个解密函数,再转换格式,然后交给另一个库消费。它像一个黑盒,数据出来就直接触发了事件,中间加工环节很别扭。
实操心得: 我曾在一个项目里用
EventSource接收日志流,需要实时高亮显示错误关键词。由于EventSource吐出来的是纯文本,我不得不在onmessage回调里用正则表达式匹配和替换 DOM 元素,当日志量巨大时,UI 线程频繁操作 DOM 导致页面卡顿。这就是“只读文本流”和“缺乏中间处理能力”带来的典型问题。如果当时能用上可转换的流,完全可以在数据流到达 UI 之前就完成文本处理,体验会好很多。
正是这些限制,催生了我们对更强大流处理能力的需求,从而引入了 Web Streams API。
3. ReadableStream:拥抱任意数据源的基础流
Web Streams API 是现代浏览器提供的一组底层 API,用于高效地处理流式数据。ReadableStream是其中的“可读流”,代表了一个数据源。
3.1 理解 ReadableStream 的模型
你可以把ReadableStream想象成一个水桶,数据是水。这个水桶有一个“龙头”(reader),你可以打开龙头来接水(读取数据)。水桶里的水可以是任何东西:文本块、二进制数组(Uint8Array)、甚至是其他语言中的结构化对象。
与EventSource最大的不同是,ReadableStream是协议无关和格式无关的。它的数据可以来自:
- Fetch API 的响应体:
response.body就是一个ReadableStream。 - 本地文件: 通过
File或Blob对象的stream()方法获得。 - WebSocket: 可以手动将 WebSocket 的消息包装成
ReadableStream。 - 自定义生成: 用
new ReadableStream()构造函数,自己创建一个流。
3.2 核心用法:两种读取模式
模式一:使用 Reader(更底层,控制力强)
// 假设我们从 fetch 得到一个流 const response = await fetch('https://api.example.com/large-file'); const readableStream = response.body; // 这就是一个 ReadableStream // 1. 获取一个读取器(Reader),这会“锁定”这个流,其他读取器无法再读取 const reader = readableStream.getReader(); try { while (true) { // 2. 读取一块数据。`done` 为 true 表示流结束了。 const { done, value } = await reader.read(); if (done) { console.log('流读取完毕'); break; } // 3. `value` 就是读取到的数据块,可能是 Uint8Array (对于二进制流) 或字符串 console.log('收到数据块:', value); // 在这里处理数据块,比如拼接、解析、展示 } } catch (error) { console.error('读取流时发生错误:', error); } finally { // 4. 非常重要!释放锁,允许其他代码再次读取这个流。 reader.releaseLock(); }模式二:使用异步迭代(更现代,代码简洁)
const response = await fetch('https://api.example.com/ndjson-stream'); // 假设是 ndjson 流 const readableStream = response.body; // 使用 for await...of 循环来迭代流 for await (const chunk of readableStream) { // chunk 同样是数据块 console.log('收到数据块:', chunk); // 注意:这里拿到的 chunk 很可能是 Uint8Array,需要解码 const text = new TextDecoder().decode(chunk); const jsonData = JSON.parse(text); // 假设是 ndjson console.log('解析后的数据:', jsonData); }3.3 关键特性:反压(Backpressure)
这是ReadableStream(乃至整个 Web Streams API)的精髓之一,也是它比EventSource高级的地方。反压是一种流控制机制,确保消费者(读取流的一方)不会被生产者(产生流的一方)的数据淹没。
原理很简单:当消费者处理数据的速度跟不上生产者发送数据的速度时,消费者可以“慢下来”,告诉生产者“我还没准备好,请暂停发送”。在ReadableStream中,这是通过reader.read()这个异步操作自然实现的。生产者(比如底层的网络栈)会等待read()被调用后才推送下一块数据。如果消费者不调用read(),数据就会在源头被缓冲或等待。
对比EventSource:EventSource没有显式的反压。数据来了就触发事件,如果你的onmessage回调函数执行太慢,事件会堆积在队列里,最终可能导致内存问题。你需要自己用队列和标志位来模拟反压,非常麻烦。
3.4 实战:用 ReadableStream 消费 Fetch 的流式 JSON
这是一个非常实用的场景。假设你的 API 返回一个流式的 NDJSON(Newline Delimited JSON),每行是一个 JSON 对象。用EventSource处理会很别扭,但用ReadableStream就非常自然:
async function consumeNDJSONStream(url) { const response = await fetch(url); const reader = response.body .pipeThrough(new TextDecoderStream()) // 先将二进制流转换为文本流 .getReader(); let buffer = ''; try { while (true) { const { done, value } = await reader.read(); if (done) break; buffer += value; const lines = buffer.split('\n'); // 最后一行可能是不完整的,留回缓冲区 buffer = lines.pop() || ''; for (const line of lines) { if (line.trim() === '') continue; // 跳过空行 try { const data = JSON.parse(line); // 处理每一个完整的 JSON 对象 console.log('收到数据:', data); // 更新UI,存储到状态管理等 } catch (e) { console.error('解析 JSON 行失败:', line, e); } } } // 循环结束后,处理缓冲区可能残留的最后一行(如果服务器以换行符结束,则 buffer 为空) if (buffer.trim()) { const data = JSON.parse(buffer); console.log('最后一行数据:', data); } } finally { reader.releaseLock(); } }踩坑记录: 在上面的例子中,
TextDecoderStream是一个TransformStream(我们下一节会详细讲)。这里直接用了。但早期没有这个类的时候,我们需要手动用TextDecoder来解码每个Uint8Array块,并小心处理跨块的字符分割问题(比如一个中文字符可能被截断在两个数据块里)。这就是直接使用底层ReadableStream时需要注意的细节。TextDecoderStream帮我们完美地解决了这个问题。
ReadableStream给了我们处理任何流的能力,但它主要解决的是“读”的问题。当我们需要在“读”和“消费”之间对数据做点什么的时候,就需要TransformStream登场了。
4. TransformStream:流数据的中转加工站
如果说ReadableStream是水源,WritableStream是目的地(比如写入文件),那么TransformStream就是连接它们之间的、可以随意组装和替换的“管道处理器”。它同时实现了可读和可写接口,写入一端的数据,经过内部转换后,从可读一端出来。
4.1 核心概念与工作原理
一个TransformStream内部有一个transformer对象,这个对象至少需要实现一个transform(chunk, controller)方法。
class MyTransformer { transform(chunk, controller) { // 1. `chunk` 是写入端传入的数据 // 2. 在这里对 chunk 进行任何处理:转换、过滤、加密、压缩等 const processedChunk = this._doSomething(chunk); // 3. 通过 controller.enqueue() 将处理后的数据送入可读端 controller.enqueue(processedChunk); // 如果需要,也可以选择不 enqueue,这就实现了“过滤” } _doSomething(chunk) { // 你的转换逻辑 return chunk.toUpperCase(); // 例如,把所有文本转大写 } } const myTransformStream = new TransformStream(new MyTransformer());然后,你可以通过.pipeThrough()方法将多个流连接起来:
// 假设 sourceStream 是一个 ReadableStream const processedStream = sourceStream .pipeThrough(new TextDecoderStream()) // 第一站:二进制转文本 .pipeThrough(new MyTransformer()) // 第二站:自定义转换(如转大写) .pipeThrough(new TextEncoderStream()); // 第三站:文本转回二进制(如果需要) // processedStream 仍然是一个 ReadableStream,可以继续被消费这种“管道式”编程模型非常清晰和强大,每个TransformStream职责单一,易于测试和复用。
4.2 内置的 TransformStream
浏览器已经为我们提供了一些非常实用的内置转换流:
TextDecoderStream: 将Uint8Array(二进制)流转换为字符串流。TextEncoderStream: 将字符串流转换回Uint8Array流。CompressionStream/DecompressionStream: 用于 gzip 或 deflate 格式的压缩和解压缩流数据。ByteLengthQueuingStrategy和CountQueuingStrategy: 这些是用于控制流队列策略的,通常不直接实例化,但在创建自定义流时有用。
4.3 实战:构建一个 SSE 到 JSON 对象的转换流
还记得EventSource只能吐文本,且难以中间处理的问题吗?现在,我们可以用ReadableStream和TransformStream来构建一个更强大的“增强版 EventSource”。
假设我们有一个 SSE 端点,但我们想用流的方式处理,并且直接得到解析好的 JSON 对象。
步骤1:用 Fetch 读取 SSE 流EventSource底层也是 HTTP 请求,我们可以用fetch来发起同样的请求,但获得一个ReadableStream。
async function createEnhancedEventSource(url) { const response = await fetch(url, { headers: { 'Accept': 'text/event-stream', // 告诉服务器我们需要 SSE }, }); if (!response.ok || !response.body) { throw new Error(`SSE 连接失败: ${response.status}`); } // response.body 是一个 ReadableStream (二进制流) return response.body; }步骤2:创建 SSE 解析转换流这是核心。我们需要解析data:、event:、id:等 SSE 格式。
class SSETransformer { constructor() { this.buffer = ''; this.eventType = 'message'; this.lastEventId = ''; } transform(chunk, controller) { // chunk 现在是字符串(因为已经过 TextDecoderStream) this.buffer += chunk; const lines = this.buffer.split('\n'); this.buffer = lines.pop() || ''; // 剩余的不完整行放回缓冲区 let currentEvent = { type: this.eventType, id: this.lastEventId, data: '' }; for (const line of lines) { if (line.startsWith('event:')) { currentEvent.type = line.substring(6).trim(); } else if (line.startsWith('data:')) { currentEvent.data += (currentEvent.data ? '\n' : '') + line.substring(5).trim(); } else if (line.startsWith('id:')) { currentEvent.id = line.substring(3).trim(); this.lastEventId = currentEvent.id; } else if (line.startsWith('retry:')) { // 可以处理重连时间,这里略过 } else if (line.trim() === '') { // 空行表示一个事件结束 if (currentEvent.data !== '') { // 将解析好的事件对象送入下游 controller.enqueue(currentEvent); } // 重置当前事件,为下一个事件做准备 currentEvent = { type: this.eventType, id: this.lastEventId, data: '' }; } // 忽略以冒号开头的注释行 } // 处理缓冲区结束后可能残留的最后一个事件(如果末尾有空行,则已处理) if (currentEvent.data !== '' && this.buffer === '') { controller.enqueue(currentEvent); } } flush(controller) { // 流结束时,如果缓冲区还有数据,尝试触发最后一个事件 if (this.buffer.trim()) { // 这里简化处理,实际可能需要更严谨的解析 controller.enqueue({ type: this.eventType, data: this.buffer.trim(), id: this.lastEventId }); } } }步骤3:组装完整的处理管道现在,我们把它们连起来:
async function connectToSSE(url) { try { const rawStream = await createEnhancedEventSource(url); // 步骤1:获取原始二进制流 const jsonStream = rawStream .pipeThrough(new TextDecoderStream()) // 二进制 -> 文本 .pipeThrough(new TransformStream(new SSETransformer())) // 文本 -> SSE事件对象 .pipeThrough(new TransformStream({ transform(event, controller) { // 将 SSE 事件对象转换为业务需要的格式,例如解析 JSON data try { if (event.data) { const parsedData = JSON.parse(event.data); controller.enqueue({ type: event.type, id: event.id, data: parsedData, // 这里是解析后的 JSON raw: event.data // 保留原始数据以备不时之需 }); } } catch (e) { console.warn('解析 JSON 失败:', event.data, e); // 可以选择将错误事件传递下去,或者忽略 controller.enqueue({ ...event, parseError: e }); } } })); // 现在,jsonStream 是一个 ReadableStream,它产出的是已经解析好的、带类型和ID的事件对象 const reader = jsonStream.getReader(); while (true) { const { done, value } = await reader.read(); if (done) { console.log('SSE 流正常结束'); break; } // value 现在是一个结构清晰的 JavaScript 对象 console.log(`[${value.type}] ID:${value.id}`, value.data); // 根据 value.type 分发到不同的处理函数,就像 EventSource 的 addEventListener 一样 handleEvent(value.type, value.data); } } catch (error) { console.error('连接或处理 SSE 失败:', error); // 这里可以实现自己的重连逻辑,比 EventSource 的自动重连更灵活 setTimeout(() => connectToSSE(url), 5000); } }这个方案的优势:
- 格式灵活: 最终得到的是 JavaScript 对象,
data字段已经是解析好的 JSON,无需在业务代码中再调用JSON.parse。 - 强大的中间处理能力: 在管道中,我们可以轻松插入其他
TransformStream。例如,可以插入一个流来过滤某些类型的事件,或者对数据进行加密/解密,或者将多个 SSE 流合并。 - 完整的流控制: 受益于
ReadableStream的反压机制,如果下游处理慢,上游的 fetch 请求也会慢下来,不会压垮客户端。 - 更好的错误处理和重连控制: 我们可以实现比
EventSource更精细的重连策略(比如指数退避、根据错误类型决定是否重连)。
深度解析: 你可能注意到,我们失去了
EventSource的自动重连和Last-Event-ID机制。这是因为我们接管了底层的 HTTP 请求。要实现重连,我们需要在catch块中手动重试。要实现断点续传,我们需要在连接断开时,记录最后一个收到的event.id,并在下一次连接的请求头中手动设置Last-Event-ID。这增加了代码量,但也带来了极大的灵活性。例如,你可以只在网络错误时重连,而在服务器返回 4xx 错误时停止。你也可以实现更复杂的重连间隔算法。
5. 进化之路:如何为你的项目选择?
现在,我们清晰地看到了从EventSource到ReadableStream再到TransformStream的进化路径。它们不是简单的替代关系,而是提供了不同层次的抽象和能力。
选择EventSource当:
- 你的需求非常简单:只是接收服务器推送的文本消息。
- 你需要开箱即用的自动重连和断点续传。
- 你的目标浏览器环境可能非常古老(虽然现在大部分现代浏览器都支持 Streams API,但
EventSource的历史更久)。 - 你不想在前端处理任何流解析逻辑,希望保持客户端代码极简。
选择ReadableStream(配合fetch)当:
- 你需要处理非 SSE 的流式数据,例如大文件下载、流式 API 响应(如 NDJSON)。
- 你需要二进制数据,或者想自己控制数据块的读取节奏(反压)。
- 你正在使用
fetch,并且想充分利用响应体的流式特性来优化内存使用(例如,边下载边处理一个大文件,而不是等全部下载完)。
选择组合使用ReadableStream和TransformStream当:
- 你需要一个“增强版 EventSource”,在传输过程中对数据进行复杂的转换、过滤、合并。
- 你正在构建一个需要处理多种流式数据源,并进行统一处理的中间件或库。
- 你对性能有极致要求,希望实现从网络到业务逻辑的无缝、高效流水线,避免中间不必要的拷贝和停顿。
一个更现代的架构思考:在现代前端架构中,特别是基于状态管理(如 Redux、Vuex)或响应式编程(如 RxJS)的应用中,流的思想可以很好地融入。你可以创建一个“数据流服务”,它使用ReadableStream+TransformStream从服务器获取并处理数据,然后将处理好的数据对象推送到一个 Observable 或直接 dispatch 到状态库中。这样,UI 组件只需要订阅状态变化,完全不用关心数据是怎么来的、怎么解析的,实现了极佳的关注点分离。
从EventSource到ReadableStream再到TransformStream,这条进化路径反映了前端对数据处理的掌控力从“应用层”深入到“传输层”乃至“字节流层”的过程。理解它们,不仅能帮你解决眼前的具体问题,更能让你在面对未来更复杂的实时数据场景时,拥有从底层构建解决方案的能力。下次当你需要处理源源不断的数据时,不妨先想想:这次,我该用哪一级的“武器”?