【问题标题】:node.js - create a new ReadStream for a new file, when that file reaches a certain sizenode.js - 当文件达到一定大小时,为新文件创建一个新的 ReadStream
【发布时间】:2016-08-02 17:55:27
【问题描述】:

继续Node - how can i pipe to a new READABLE stream?

我正在尝试使用fs.watchfs.stat 为我的实时编码MP3 文件创建一个新的ReadStream,当它达到一定大小(基本上是预缓冲)时。

它可以工作,但是一旦ReadStream 启动,我不知道如何退出观察程序并保持流运行。

我尝试了如下的承诺,但从未解决,所以streamEncodedFile 被反复调用:

var watcher = fs.watch(mp3RecordingFile);

watcher.on('change', (event, path) => {

  fs.stat(mp3RecordingFile, function (err, stats) {

    if (stats.size > 75533) {

        new Promise(function(resolve, reject) {
                streamEncodedFile(); 
        })
        .then(function(result) {
                watcher.close();
                console.log('watcher closed');
        });

    }

  });
});

function streamEncodedFile() {

  var mp3File = fs.createReadStream(mp3RecordingFile);

            mp3File.on('data', function(buffer){
                io.sockets.emit('audio', { buffer: buffer });
            });

}

我的另一个可悲的尝试是尝试仅以特定文件大小启动流:

watcher.on('change', (event, path) => {

  fs.stat(mp3RecordingFile, function (err, stats) {

    console.log(stats.size);

    if (stats.size > 75533 && stats.size < 75535) {
                streamEncodedFile(); 
    } else if (stats.size > 75535) {
                watcher.close(); 
    } 

  });
});

【问题讨论】:

  • var watcher = fs.watch(mp3RecordingFile); watcher.on('change', (event, path) =&gt; { fs.stat(mp3RecordingFile, function(err, stats) { if (stats.size &gt; 75533) { streamEncodedFile(); } }); }); function streamEncodedFile() { var mp3File = fs.createReadStream(mp3RecordingFile); mp3File.on('data', function(buffer) { io.sockets.emit('audio', { buffer: buffer }); }); watcher.close(); }
  • 谢谢。我应该说我试过了,但是 ReadStream 结束了。我相信这就是为什么github.com/jasontbradshaw/tailing-stream

标签: javascript node.js stream


【解决方案1】:

试试这个解决方案,缓冲和写入文件。

const Writable = require('stream').Writable;
const fs = require('fs');

let mp3File = fs.createWriteStream('path/to/file.mp3');

var buffer = new Buffer([]);
//in bytes
const CHUNK_SIZE = 102400; //100kb

//Proxy for emitting and writing to file
const myWritable = new Writable({
  write(chunk, encoding, callback) {
    buffer = Buffer.concat([buffer, chunk]);
    if(buffer.length >= CHUNK_SIZE) {
       mp3File.write(buffer);
       io.sockets.emit('audio', { buffer: buffer});
       buffer = new Buffer([]);
    }

    callback();
  }
});

myWritable.on('finish', () => {
   //emit final part if there is data to emit
   if(buffer.length) {
       //write final chunk and close fd
       mp3File.end(buffer);
       io.sockets.emit('audio', { buffer: buffer});
   }
});


inbound_stream.pipe(encoder).pipe(myWritable);

【讨论】:

  • 迫不及待想试试这个。再次感谢!
  • 再次 @Nazar Sakharenko。我不明白你为什么要用+= 附加到缓冲区。块是 ArrayBuffer 对象,需要保持这种方式才能在另一端解码。我们不应该做buffer.push(chunk)吗?然后我需要以某种方式迭代数组以开始从中发送原始的 ArrayBuffers。除非我错过了什么。谢谢!
  • 不,我刚刚忘记了你在操作缓冲区。答案已更新。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-02-04
  • 2011-12-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多