Parallel stream huge line delimited json file in nodejs

Viewed 410

I am reading a file with 350M lines using createReadstream and transforming each line and writing it back as line delimited file. Below is the code which I am using to do it.

var fs = require("fs");
var args = process.argv.slice(2);
var split = require("split")
fs.createReadStream(args[0])
    .pipe(split(JSON.parse))
    .on('data', function(obj) {
        <data trasformation operation>
    })
    .on('error', function(err) {
    })

To red 350M lines it takes 40 minutes and it only uses one CPU core while doing it. I have 16 CPU cores. How can I make this line reading process to run parallel so that alteast 10 cores are utilized and the entire operation finishes in less time.

I tried using this module - https://www.npmjs.com/package/parallel-transform. But when I checked in htop, it was still single CPU which was doing the operation.

var stream = transform(10, {
    objectMode: true
}, function(data, callback) {
    <data trasformation operation>
    callback(null, data);
});

fs.createReadStream(args[0])
    .pipe(stream)
    .pipe(process.stdout);

What is the better way to parallel read file while streaming?

1 Answers

You can try scramjet - I would be happy to find someone with a strong multi-threaded use case for setting up proper tests around this.

Your code would look something like this:

var fs = require("fs");
var {StringStream} = require("scramjet");
var args = process.argv.slice(2);

let i = 0;
let threads = os.cpus().length; // you may want to check this out

StringStream.from(fs.createReadStream(args[0]))
    .lines() // it's better to deserialize this in the threads
    .separate(() => i = ++i % threads)
    .cluster(stream => stream // these will happen in the thread
        .JSONParse()
        .map(yourProcessingFunc) // this can be async as well
    )
    .mux() // if the function above returns something you'll get
           // a stream of results
    .run() // this executes the whole workflow.
    .catch(errorHandler)

You can use a better affinity function for separate, see the docs here where based on the data you can direct the data to specific workers. If you encounter any issues, please create a repo and let's see how to fix those.

Related