【问题标题】:How to properly handle back-pressure during `Transform#flush`如何在“Transform#flush”期间正确处理背压
【发布时间】:2019-01-05 09:21:06
【问题描述】:

在 Transform 的 _flush 方法的实现中处理背压的正确方法是什么?换句话说,如果.push() 在刷新时返回 false,是否有任何机制可以正确处理来自下游的背压?

文档规定一旦 .push() 返回 false 就停止推送,但是当下游想要恢复读取时,Transform 没有办法监听,只能覆盖 this.read;但这会是什么样子?这样做有什么危险吗?

这是一个你可以玩的工作示例。

const stream = require('stream');

// a string large enough to overflow the buffer
const S_OVERFLOW = '-'.repeat((new stream.Readable()).readableHighWaterMark+1);


class example extends stream.Transform {
    constructor() {
        super({
            writableObjectMode: true,
        });

        // some internal queue that will be emptied once writable side ends
        Object.assign(this, {
            internal_queue: [],
        });
    }

    _transform(g_chunk, s_encoding, fk_transform) {
        // store chunk in internal queue
        this.internal_queue.push(g_chunk);

        // done with transform (no writes)
        fk_transform();
    }

    _flush(fk_flush) {
        console.warn('starting to flush');

        // now that writable side has ended, flush internal queue
        this.resumeFlush(fk_flush);
    }

    resumeFlush(fk_flush) {
        let a_queue = this.internal_queue;

        // still data left in internal queue
        while(a_queue.length) {
            // remove an item from queue
            a_queue.pop();

            // intentionally overflow buffer
            if(!this.push(S_OVERFLOW)) {
                //
                // WHAT TO DO HERE?
                //

                // go asynchronous
                return;
            }
        }

        console.warn('finished flush');

        // callback
        fk_flush();
    }
}


// instantiate transform
let ds_transform = new example();

// pipe to stdout
ds_transform.pipe(process.stdout);

// write some data (needs to happen twice)
ds_transform.write({
    item: 0,
});

ds_transform.write({
    item: 1,
});

// end stream
ds_transform.end();

将 stdout 连接到 /dev/null 以便 stderr 仍然打印到控制台:

$ node transform.js > /dev/null
starting to flush

【问题讨论】:

    标签: node.js node-streams


    【解决方案1】:

    这里真正的问题是您应该使用 Duplex 而不是 Transform。由于对_transform 的每次调用实际上是在缓冲数据而不是对其应用一些(a/)同步转换,因此这种类型的实现更适合作为双工,即调用_write() 缓冲数据,调用@987654323 @ 开始推送,直到检测到背压。

    const stream = require('stream');
    
    // a string large enough to overflow the buffer
    const S_OVERFLOW = '-'.repeat((new stream.Readable()).readableHighWaterMark+1);
    
    
    class example extends stream.Duplex {
        constructor() {
            super({
                writableObjectMode: true,
            });
    
            // some internal queue that will be emptied once writable side ends
            Object.assign(this, {
                internal_queue: [],
            });
        }
    
        _write(g_chunk, s_encoding, fk_write) {
            // store chunk in internal queue
            this.internal_queue.push(g_chunk);
    
            // done with transform (no writes)
            fk_write();
        }
    
        _read() {
            console.warn('called _read()');
            let a_queue = this.internal_queue;
    
            // still data left in internal queue
            while(a_queue.length) {
                // remove an item from queue
                a_queue.pop();
    
                // intentionally overflow buffer
                if(!this.push(S_OVERFLOW)) {
                    // go asynchronous
                    return;
                }
            }
    
            console.warn('finished reading');
    
            // nothing more to read
            this.push(null);
        }
    }
    
    
    // instantiate transform
    let ds_transform = new example();
    
    // pipe to stdout
    ds_transform.pipe(process.stdout);
    
    // write some data (needs to happen twice)
    ds_transform.write({
        item: 0,
    });
    
    ds_transform.write({
        item: 1,
    });
    
    // end stream
    ds_transform.end();
    

    然后你会得到:

    $ node duplex.js > /dev/null
    called _read()
    called _read()
    called _read()
    finished reading
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-31
      • 1970-01-01
      • 1970-01-01
      • 2023-01-20
      • 2014-01-13
      • 1970-01-01
      相关资源
      最近更新 更多