Node.js如何创建可追加可读流处理WebSocket分块Base64数据
Node.js WebSocket分块Base64转可追加可读流问题解决
问题说明
- 通过WebSocket接收前端发来的分块Base64编码字符串,需要将其转换为可持续追加数据的可读流,而非每次创建新流
- 当前实现中出现错误:
Error [ERR_METHOD_NOT_IMPLEMENTED]: The _read() method is not implemented
解决方案
1. 实现自定义可追加的Readable流
Node.js的Readable流要求必须实现_read()方法,否则会抛出上述错误。我们可以继承Readable类,实现一个支持追加数据的自定义流:
const { Readable } = require('stream'); class AppendableReadableStream extends Readable { constructor(options = {}) { super(options); this.isClosed = false; } // 必须实现的_read方法,空实现即可满足基础流机制要求 _read() {} // 追加Base64数据到流,自动解码为Buffer append(base64Data) { if (this.isClosed) { throw new Error('流已关闭,无法追加数据'); } const buffer = Buffer.from(base64Data, 'base64'); this.push(buffer); } // 结束流,告知消费者无更多数据 close() { this.isClosed = true; this.push(null); } }
2. 在WebSocket连接中复用流实例
在新的WebSocket连接建立时初始化流,后续收到分块数据直接追加到同一个流,直到前端发送结束信号或连接断开:
// 每个WebSocket连接对应一个流(根据业务需求调整作用域) let audioStream; socket.on('connection', (clientSocket) => { // 新连接建立时创建可追加流 audioStream = new AppendableReadableStream(); clientSocket.on('message', async (res) => { try { const payload = JSON.parse(res); const payloadType = payload.type; const data = payload.data; // 处理分块数据 if (payloadType === 'audio_chunk') { audioStream.append(data); } // 处理结束信号 else if (payloadType === 'end') { audioStream.close(); } } catch (err) { console.error('消息处理出错:', err); audioStream?.close(); } }); // 连接断开时自动关闭流,避免内存泄漏 clientSocket.on('close', () => { audioStream?.close(); }); });
3. 适配原有流消费逻辑
直接使用自定义流替换原有micStream,即可正常消费持续追加的数据流:
const SAMPLE_RATE = 16000; // 根据你的实际采样率调整 function encodePCMChunk(chunk) { // 保留原有的PCM编码逻辑 return chunk; } const getAudioStream = async function* () { for await (const chunk of audioStream) { if (chunk.length <= SAMPLE_RATE) { yield { AudioEvent: { AudioChunk: encodePCMChunk(chunk), }, }; } } };
内容的提问来源于stack exchange,提问作者james Makinde
相关产品推荐
相关产品推荐

