【问题标题】:Pausing cassandra stream for async operations为异步操作暂停 cassandra 流
【发布时间】:2020-01-18 07:08:18
【问题描述】:

我想在处理下一行之前暂停我的 cassandra 流以进行一些异步操作。

每一行都在一个可读的事件侦听器中接收。我试过使用 stream.pause 但它实际上并没有暂停流。 我也在“数据”事件侦听器中尝试过同样的方法,但这也不起作用。 也许会非常感谢您的见解和解决方案。 这是我的代码。 在可读和“等待”中使用 async 使用 await 实际上并不能阻止下一行在异步函数完成之前出现。

function start() {
let stream = client.stream('SELECT * FROM table');
stream
    .on('end', function () {
        console.log(`Ended at ${Date.now()}`);
    })
    .on('error', function (err) {
        console.error(err);
    })
    .on('readable', function () {
        let row = this.read();
        asyncFunctionNeedTowaitForthisBeforeNextRow()
    })
}

//下面的不行

function start() {
let stream = client.stream('SELECT * FROM table');
stream
    .on('end', function () {
        console.log(`Ended at ${Date.now()}`);
    })
    .on('error', function (err) {
        console.error(err);
    })
    .on('readable', async function () {
        let row = this.read();
        stream.pause();
        await asyncFunctionNeedTowaitForthisBeforeNextRow();
        stream.resume();
    })
 }

【问题讨论】:

    标签: node.js cassandra stream cassandra-driver express-cassandra


    【解决方案1】:

    stream.pause() 不起作用的原因是因为readable 事件触发了多次,所以再次调用了同一个异步函数。 data 事件也是如此。

    我建议使用自定义 writable stream 来正确处理所有这些异步内容。

    可写流如下所示:

    const {Writable} = require('stream');
    
    const myWritable = new Writable({
      async write(chunk, encoding, callback) {
        let row = chunk.toString();
        await asyncFunctionNeedTowaitForthisBeforeNextRow();
        callback(); // Write completed successfully
      }
    })
    

    然后调整您的代码以使用此可写:

    function start() {
      let stream = client.stream('SELECT * FROM table');
      stream.pipe(myWritable);
      stream
        .on('end', function () {
          console.log(`Ended at ${Date.now()}`);
        })
        .on('error', function (err) {
          console.error(err);
        })
    }
    

    【讨论】:

      【解决方案2】:

      请注意,即使您将事件 'readable' 的处理程序声明为异步函数,调用者也不会等待返回的承诺完成,因为 Stream 期望事件处理程序正常执行。

      解决方案可能是:

      stream
        .on('end', () => {})
        .on('error', () => {})
        .on('data', row => {
          stream.pause();
          doSomethingAsync(row).then(() => stream.resume());
        });
      

      请注意,理想情况下,您应该在执行异步操作时利用并行性,因此最好每次读取几行然后暂停。

      【讨论】:

        猜你喜欢
        • 2018-01-03
        • 1970-01-01
        • 2021-12-27
        • 1970-01-01
        • 2022-07-16
        • 2021-08-24
        • 2023-03-05
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多