Node.Js async iterator over stream pipeline

Viewed 2669

I have the following pipeline:

readFile > parseCSV > otherProcess

The readFile is the standard Node.Js createReadStream, while the parseCSV is a Node.js transform stream (module link).

I want to iterate through a csv file line by line and handle a single line at the time. Therefore, streams and async iterator are a perfect match.

I have the following code which is working properly:

async function* readByLine(path, opt) {
  const readFileStream = fs.createReadStream(path);
  const csvParser = parse(opt);
  const parser = readFileStream.pipe(csvParser);
  for await (const record of parser) {
    yield record;
  }
}

I'm quite new to Node.Js streams, but I've read from many sources that the module stream.pipeline is preferred to the .pipe method of read streams.

How can I change the code above in order to use the stream.pipeline (actually the promise version got from util.promisify(pipeline)) and yielding one line at the time?

2 Answers

You should actually be able to just pass both the fs-stream and the parser-stream to pipeline() and use your async iterator on the parser-stream:

const fs = require('fs');
const parse = require('csv-parse');
const stream = require('stream')
const util = require('util');
const pipeline = util.promisify(stream.pipeline);

async function* readByLine(path, opt) {
    const readFileStream = fs.createReadStream(path);
    const csvParser = parse(opt);
    await pipeline(readFileStream, csvParser);
    for await (const record of csvParser) {
        yield record;
    }
}

Adding to @eol's answer, I would recommend storing the promise and awaiting it after the async iteration.

const fs = require('fs');
const parse = require('csv-parse');
const stream = require('stream');

async function* readByLine(path, opt) {
    const readFileStream = fs.createReadStream(path);
    const csvParser = parse(opt);
    const promise = stream.promises.pipeline(readFileStream, csvParser);
    for await (const record of csvParser) {
        yield record;
    }
    await promise;
}

By calling await pipeline(...) before the loop, it will consume the whole stream before you can iterate from whatever is left in the buffer, which works by accident on small streams but is likely to break on larger (or infinite/lazy) streams.

The callback equivalent might make it more clear what's going on depending on where we await.

// await before iterating
stream.pipeline(a, b, err => {
  if (err) return callback(err)

  for await (const record of b) {
    // process record
  }

  callback()
}

// await after iterating
for await (const record of stream.pipeline(a, b, callback)) {
  // process record
}
Related