How to clone a readable stream

Viewed 2622

I have a stream i am trying to submit the same stream to two different destinations. The first destination is to AWS S3, the second destination is to some other backend via a http request.

const document = fs.createReadStream(process.cwd() + "/test/resources/" + "id/document.jpg");

const s3Response = await submitToS3(document);

const backendResponse = await submitToBackend(document);

From what i understand a stream can only be read once. How can i send the same stream to two different destinations.

I thought about cloning the stream but simply creating a new variable and assigning the stream to that variable does not work.

3 Answers

you can check this npm module out: https://www.npmjs.com/package/readable-stream-clone

npm install readable-stream-clone
const fs = require("fs");
const ReadableStreamClone = require("readable-stream-clone");
 
const readStream = fs.createReadStream('text.txt');
 
const readStream1 = new ReadableStreamClone(readStream);
const readStream2 = new ReadableStreamClone(readStream);
 
const writeStream1 = fs.createWriteStream('sample1.txt');
const writeStream2 = fs.createWriteStream('sample2.txt');
 
readStream1.pipe(writeStream1)
readStream2.pipe(writeStream2)

Raghavendra's answer suggests a good potential direction. You can combine multiple pipes from this answer with S3 pipe implementation from this answer.

For the submitToBackend part, not sure exactly what your implementation looks like, but assuming you can pipe to an HTTP request of some sort...

Example:

var fs = require("fs");
const request = require("request");
const AWS = require("aws-sdk");
const s3 = new AWS.S3();

const rs = fs.createReadStream(process.cwd() + "/test/resources/" + "id/document.jpg");

function uploadFromStream(s3) {
  var pass = new stream.PassThrough();

  var params = {Bucket: BUCKET, Key: KEY, Body: pass};
  s3.upload(params, function(err, data) {
    console.log(err, data);
  });

  return pass;
}

rs.pipe(uploadFromStream(s3));

// Just guessing at your submitToBackend implementation:
const backendWs = request.post("http://example.com/docs");

// However it works, if you can get to a stream.Writable, you can now pipe the same stream.Readable:
rs.pipe(backendWs);
var fs = require("fs");
var ReadableStreamClone = require("readable_stream");

var readStream = fs.createReadStream('text1.txt');

var readStream1 = new ReadableStreamClone(readStream);
var readStream2 = new ReadableStreamClone(readStream);

var writeStream1 = fs.createWriteStream('testsample1.txt');
var writeStream2 = fs.createWriteStream('testsample2.txt');

readStream1.pipe(writeStream1)
readStream2.pipe(writeStream2)
Related