【发布时间】: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