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

如何用Python或NodeJS分块发送数据?含Node.js pipe()用法疑问

分块发送音频数据的实现方案(Python/Node.js)

Python 实现分块发送

本地文件分块传输

先根据音频比特率计算10秒内容对应的字节数(比如128kbps的音频:128kbps = 16KB/s,10秒即160KB),再分块读取文件并发送,可搭配简单的块标识保证传输可靠性:

import socket
import os

# 按音频比特率计算10秒对应的字节量
bitrate_kbps = 128
bytes_per_second = (bitrate_kbps * 1024) // 8
chunk_size = bytes_per_second * 10

# 发送端:分块读取音频文件并发送
def send_audio_chunk(host, port, audio_path):
    with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
        s.connect((host, port))
        with open(audio_path, 'rb') as f:
            chunk_index = 1
            while True:
                chunk = f.read(chunk_size)
                if not chunk:
                    break
                # 发送块头(包含块长度)+ 块内容
                header = f"CHUNK:{chunk_index}|{len(chunk)}|".encode()
                s.sendall(header + chunk)
                # 等待接收端确认(可选,确保块传输完整)
                ack = s.recv(1024)
                if ack != b'ACK':
                    print(f"块{chunk_index}传输失败")
                    break
                chunk_index += 1

# 接收端:解析块头并拼接音频内容
def receive_audio_chunk(host, port, save_path):
    with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
        s.bind((host, port))
        s.listen()
        conn, addr = s.accept()
        with conn, open(save_path, 'wb') as f:
            buffer = b''
            while True:
                buffer += conn.recv(4096)
                # 解析块头分隔符
                pipe_pos = buffer.find(b'|')
                if pipe_pos == -1:
                    continue
                # 获取块序号和长度
                chunk_info = buffer[:pipe_pos].decode().split(':')
                chunk_len = int(chunk_info[1])
                # 提取块内容
                if len(buffer) >= pipe_pos + 1 + chunk_len:
                    chunk = buffer[pipe_pos+1 : pipe_pos+1+chunk_len]
                    f.write(chunk)
                    conn.sendall(b'ACK')
                    # 重置缓冲区,处理剩余数据
                    buffer = buffer[pipe_pos+1+chunk_len:]
                if not buffer:
                    break

实时音频流分块

如果是实时生成的音频,可通过生成器逐块输出音频数据,再重复上述发送逻辑。


Node.js 实现分块发送

手动分块传输

逻辑和Python一致,先计算10秒字节量,再分块读写:

const fs = require('fs');
const net = require('net');

const bitrateKbps = 128;
const bytesPerSecond = (bitrateKbps * 1024) / 8;
const chunkSize = bytesPerSecond * 10;

// 发送端
function sendAudioChunk(host, port, audioPath) {
  const client = net.connect(port, host, () => {
    const readStream = fs.createReadStream(audioPath, { highWaterMark: chunkSize });
    let chunkIndex = 1;
    readStream.on('data', (chunk) => {
      const header = `CHUNK:${chunkIndex}|${chunk.length}|`;
      client.write(Buffer.from(header));
      client.write(chunk);
      // 等待确认
      client.once('data', (ack) => {
        if (ack.toString() !== 'ACK') readStream.destroy();
        chunkIndex++;
      });
    });
    readStream.on('end', () => client.end());
  });
}

// 接收端
function receiveAudioChunk(host, port, savePath) {
  const server = net.createServer((socket) => {
    const writeStream = fs.createWriteStream(savePath);
    let buffer = Buffer.alloc(0);
    let expectedChunkLen = null;

    socket.on('data', (data) => {
      buffer = Buffer.concat([buffer, data]);
      if (!expectedChunkLen) {
        const pipePos = buffer.indexOf('|');
        if (pipePos === -1) return;
        const chunkLen = parseInt(buffer.slice(0, pipePos).toString().split(':')[1]);
        expectedChunkLen = chunkLen;
        buffer = buffer.slice(pipePos + 1);
      }
      if (buffer.length >= expectedChunkLen) {
        const chunk = buffer.slice(0, expectedChunkLen);
        writeStream.write(chunk);
        socket.write('ACK');
        buffer = buffer.slice(expectedChunkLen);
        expectedChunkLen = null;
      }
    });
    socket.on('end', () => writeStream.end());
  });
  server.listen(port, host);
}

用 pipe() 实现流式分块

你提到的pipe()是Node.js流的核心方法,它能自动处理数据分块、背压(防止发送过快导致内存溢出),完全适配Spotify这类流媒体的分段传输逻辑:

Spotify的传输本质是HTTP分块传输编码(Chunked Transfer Encoding),将音频切割为固定时长片段后流式发送,用pipe()可以简化整个流程:

// 服务端(Express + 流式输出)
const express = require('express');
const fs = require('fs');
const app = express();

const bitrateKbps = 128;
const bytesPerSecond = (bitrateKbps * 1024) / 8;
const chunkSize = bytesPerSecond * 10; // 对应10秒音频

app.get('/stream-audio', (req, res) => {
  const audioPath = './target-audio.mp3';
  const stat = fs.statSync(audioPath);
  
  // 设置响应头,启用分块传输
  res.setHeader('Content-Type', 'audio/mpeg');
  res.setHeader('Content-Length', stat.size);
  res.setHeader('Transfer-Encoding', 'chunked');

  // 创建读取流,设置分块大小为10秒对应字节数
  const readStream = fs.createReadStream(audioPath, { highWaterMark: chunkSize });
  // 用pipe自动将分块数据写入响应流
  readStream.pipe(res);

  readStream.on('error', () => res.status(500).end());
});

app.listen(3000);

客户端接收流式音频:

const http = require('http');
const fs = require('fs');

const writeStream = fs.createWriteStream('./received-audio.mp3');

http.get('http://localhost:3000/stream-audio', (res) => {
  // 直接pipe到文件流,自动接收分块数据
  res.pipe(writeStream);
  res.on('end', () => console.log('音频接收完成'));
});

pipe()会自动将读取流的内容按设置的highWaterMark分块,写入到可写流(这里是HTTP响应流/文件流),无需手动处理分块逻辑,同时自动平衡读写速度,避免内存堆积。


内容的提问来源于stack exchange,提问作者Alexis Mendez

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 06:11:21