promise中处理可读流需转换为promise可消费形式,常用方法包括转为buffer/字符串、for await...of迭代、pipeline汇总或手动拼接chunks,兼顾错误、超时与内存控制。

在 Promise 中处理可读流(Readable Stream)不能直接用 await 或 .then() 等待流“完成”,因为流是逐步推送数据的异步迭代器,不是一次性返回值的 Promise。正确方式是将流转换为 Promise 可消费的形式——最常用的是转成字符串、Buffer 或数组,或使用 for await...of 配合封装逻辑。
用 pipeline + Promise 汇总流到 Buffer(Node.js)
适用于 Node.js 环境中的 fs.createReadStream、HTTP 响应流等。推荐使用 stream.pipeline 配合 Uint8Array 收集器,避免手动监听 data 事件出错:
- 导入
stream/promises(Node.js ≥ 16.14)或util.promisify(pipeline) - 创建一个
Writable流(如new PassThrough())或直接用stream.Readable.toWeb()(较新环境)不适用,应选传统收集方式 - 更稳妥做法:用
new Promise+collect chunks手动拼接
示例(兼容性好):
function streamToBuffer(readable) {
const chunks = [];
return new Promise((resolve, reject) => {
readable.on('data', chunk => chunks.push(chunk));
readable.on('error', reject);
readable.on('end', () => resolve(Buffer.concat(chunks)));
});
}
<p>// 使用
const fs = require('fs');
const stream = fs.createReadStream('file.txt');
streamToBuffer(stream).then(buf => console.log(buf.toString()));
</p>
用 for await...of 消费流(现代 Node.js / 浏览器 ReadableStream)
Node.js ≥ 12.10 和现代浏览器支持可迭代流(ReadableStream 实现 [Symbol.asyncIterator])。这是最自然、可中断、内存友好的方式:
- 确保流是
AsyncIterable(Node.js 的fs.createReadStream默认不是,需用stream.Readable.toWeb()或包装;而fetch().body是原生可迭代流) - 用
for await (const chunk of stream)逐块处理,适合大文件、不需要全量加载的场景 - 可在循环中
break或return提前退出
示例(浏览器 fetch 流):
async function readResponseBody(response) {
const reader = response.body.getReader();
let result = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
result += new TextDecoder().decode(value, { stream: true });
}
return result;
}
// 或更简洁(若 body 支持 async iteration):
async function readAsText(stream) {
const decoder = new TextDecoder();
let text = '';
for await (const chunk of stream) {
text += decoder.decode(chunk, { stream: true });
}
return text;
}
封装成 Promise 工具函数(兼顾错误、超时、限流)
生产环境建议封装健壮函数,处理常见边界情况:
- 自动监听
'error'并 reject - 添加超时控制(防止流挂起)
- 限制最大缓冲大小,防内存溢出(尤其处理上传或网络流)
- 支持指定编码(如 'utf8')或返回 Uint8Array
简化版带超时:
function streamToText(readable, options = {}) {
const { encoding = 'utf8', timeout = 30_000, maxSize = 10 * 1024 * 1024 } = options;
let totalSize = 0;
const chunks = [];
<p>return Promise.race([
new Promise((resolve, reject) => {
readable.on('data', chunk => {
totalSize += chunk.length;
if (totalSize > maxSize) {
reject(new Error('Stream exceeds max size'));
readable.destroy();
return;
}
chunks.push(chunk);
});
readable.on('error', reject);
readable.on('end', () => resolve(Buffer.concat(chunks).toString(encoding)));
}),
new Promise((_, reject) => {
setTimeout(() => reject(new Error('Stream timeout')), timeout);
})
]);
}
</p>
浏览器中处理 ReadableStream(Fetch API)
浏览器原生 Response.body 是 ReadableStream,可直接用于 for await 或转为文本/数组缓冲:
-
response.text()、response.json()、response.arrayBuffer()—— 这些方法本身返回 Promise,内部已处理流,最简单推荐 - 需要自定义解析(如 CSV 行解析、JSON 分块)才需手动迭代
- 注意:调用过
text()后,body会被锁住,不能再读取
示例:
fetch('/data.json')
.then(res => res.json()) // 内部已 consume stream,返回 Promise
.then(data => console.log(data));
<p>// 自定义处理(逐行读取 NDJSON)
async function readNDJSON(stream) {
const reader = stream.getReader();
const decoder = new TextDecoder();
let buffer = '';
const results = [];
while (true) {
const { done, value } = await reader.read();
if (done && !buffer) break;
if (value) buffer += decoder.decode(value);
let i;
while ((i = buffer.indexOf('\n')) >= 0) {
const line = buffer.slice(0, i).trim();
if (line) results.push(JSON.parse(line));
buffer = buffer.slice(i + 1);
}
}
return results;
}
</p>Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











