【问题标题】:How to work properly with async operations in nodejs streams如何在 nodejs 流中正常使用异步操作
【发布时间】:2017-11-24 07:41:08
【问题描述】:

最近我开始思考——我是否以正确的方式处理 Nodejs 流中的异步操作?因此,我只想确保确实如此。

class AsyncTransform extends Transform {
   constructor() { 
      super({objectMode: true});
   }

   public async _transform(chunk, enc, done) {
      const result = await someAsyncStuff();
      result && this.push(result);
      done();
   }
}

这个小例子效果很好,但基本上,所有异步的东西都需要一些时间来执行,在大多数情况下,我想并行处理块。我会将done 放在_transform 的顶部,但由于某些原因这不是解决方案,其中一个事实是,当最后一个块将调用done 时,下一个push 将抛出错误Error: stream.push() after EOF。所以如果最后一个块已经用它的done 调用了,我们就不能推送。

为了处理这种情况,我将_flush 转换方法与块的队列计数器结合使用(当它进来时增加,push 减少)。该死的,已经有这么多字了,所以这里只是我的代码示例。

const minute = 60000;
const second = 1000;

const CRITICAL_QUEUE_POINT = 50;

export class AsyncTransform extends Transform {

    private queue: number = 0;

    constructor() {
        super({objectMode: true});
    }

    public async _transform(chunk, enc, done) {

        this.checkQueue()
            .then(() => this.init(chunk, done));
    }

    public _flush(done) {
        this._done(done, true);
    }

    private async init(chunk, done) {
        this.increaseQueueCounter();
        this._done(done);

        const user = await new UserRepository().search(chunk);

        this.decreaseQueueCounter();
        this._push(user);
    }

    /**
     * Queue
     * */
    private checkQueue(): Promise<any> {
        return new Promise((resolve) => {
            const _checkQueue = () => {
                if (this.queue >= CRITICAL_QUEUE_POINT) {
                    return setTimeout(_checkQueue, second * 10);
                }

                resolve();
            };

            _checkQueue();
        });
    }

    private increaseQueueCounter(): void {
        this.queue++;
    }

    private decreaseQueueCounter(): void {
        this.queue--;
    }

    /**
     * Transform API
     * */
    private _push(user) {
        this.push(user);
    }

    private _done(done, isFlush: boolean = false) {
        if (!isFlush) {
            return done();
        }

        if (this.queue === 0) {
            return done();
        }

        setTimeout(() => this._done(done, isFlush), second * 10);
    }

}

【问题讨论】:

    标签: javascript node.js asynchronous stream


    【解决方案1】:

    我在两年前考虑过,并正在寻找解决方案 - 某种框架。我找到了几个框架 - highland 和 event stream - 但所有框架都非常复杂,因此我决定编写一个新框架:scramjet

    现在你的代码可以这么简单:

    const {DataStream} = require('scramjet');
    
    yourDataStream.pipe(new DataStream())
        .map(async (chunk) => {
             await checkQueue();
             return new UserRepository().search(chunk);
        });
    

    或者,如果我对 checkQueue() 的理解正确,它只是将同时连接数保持在临界水平以下,那么它就更简单了:

    yourDataStream.pipe(new DataStream({maxParallel: CRITICAL_QUEUE_POINT }))
        .map(async (chunk) => new UserRepository().search(chunk));
    

    它将连接数保持在一个稳定的水平(每次有响应它都会启动一个新线程)。

    【讨论】:

    • 谢谢你,我一定会试试你的解决方案,看起来很棒!
    • 还有一个问题,我可以在 .map 之后 .pipe 另一个异步流吗?
    • 您实际上可以只使用另一个 .map - 无需像 this example 中那样使用管道。当然你仍然可以使用管道,甚至可以发送到其他流。
    猜你喜欢
    • 2021-10-16
    • 1970-01-01
    • 1970-01-01
    • 2019-01-21
    • 2016-08-24
    • 1970-01-01
    • 2015-11-13
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多