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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 23:01:23