Skip to content
This repository was archived by the owner on Jan 19, 2021. It is now read-only.
This repository was archived by the owner on Jan 19, 2021. It is now read-only.

Emits finish dispite still consuming #2

Description

@piglovesyou

The following code simply increments number asynchronously 1000 times. I assumed the parallel transform emits a finish event after all data consumption but is didn't.

const expect = 1000;
let actual = 0;
const r = createReadable(expect);

class MyTransform extends ParallelStream {
  constructor() {
    super({objectMode: true, maxParallel: 32, highWaterMark: 32});
  }

  _parallelTransform(data, enc, callback) {
    setTimeout(() => {
      console.log(++actual);
      callback();
    }, Math.random() * 100);
  }
}

const t = new MyTransform();
t.on('finish', () => {
  assert.strictEqual(actual, expect); // AssertionError [ERR_ASSERTION]: 985 === 1000
});

r.pipe(t);

function createReadable(size) {
  let i = 0;
  return new stream.Readable({
    read: function () {
      this.push(i < size ? i++ : null);
    },
    objectMode: true,
  });
}

This is not pleasant especially when we use libs such as pump that calls a callback function depending on finish event.

It seems that it emits finish after all writable data is buffered instead of consumed.

This problem occurs also in parallel-stream module.

Also, when we do readable.pipe(parallelTransform).pipe(writable), writable emits finish earlier than expected.

It is likely that we have to dig more in private methods of Node Stream.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions