【问题标题】:How to wait for each value of an observable with a promise如何等待带有承诺的可观察对象的每个值
【发布时间】:2019-07-27 11:08:30
【问题描述】:

假设我有这个 observable:

const obs = new Observable((observer) => {
    observer.next(0.25);
    observer.next(0.75);
    observer.next(new ArrayBuffer(100));
    observer.complete();
});

我怎样才能等待每个带有承诺的值?

下面的代码只会返回最后一个值(调用完成之前的值):

const value = await obs.toPromise();

但我希望能够在此过程中获得每个值。我可以这样做:

const value1 = await obs.pipe(take(1)).toPromise();
const value2 = await obs.pipe(take(2)).toPromise();

但这并不理想,因为我必须每次都增加数字,而且take(1000) 仍然会在示例中返回一些内容,即使只有 3 个值。我正在寻找类似的东西:

const value1 = await obs.pipe(next()).toPromise(); // 0.25
const value2 = await obs.pipe(next()).toPromise(); // 0.75
const value3 = await obs.pipe(next()).toPromise(); // ArrayBuffer(100)
const value4 = await obs.pipe(next()).toPromise(); // null

这更类似于生成器。

有没有办法完成这样的事情?

【问题讨论】:

  • 为什么这需要成为一个承诺?
  • @KevinB 为什么不呢?我想用 async/await 来构建我的代码。
  • 是的,但是,它看起来好像你在做任何异步的事情,所以......理论上无论有没有等待,它都会以相同的方式工作。如果您正在做的事情不是异步的,那么 async/await 将只是没有目的的额外代码。
  • 这只是一个简化的例子,请想象 observable 以异步方式发出事件。
  • @maximedupre 有一个关于converting Observable into generator 的问题,您可以尝试提供答案。虽然我仍然不确定将异步与 Observables 一起使用的大想法是否有效。

标签: javascript promise rxjs observable rxjs6


【解决方案1】:

您似乎要求的是一种将可观察对象转换为异步可迭代对象的方法,以便您可以“手动”或使用新的 for-await-of 语言功能异步迭代其值。

这是一个如何做到这一点的例子(我没有测试过这段代码,所以它可能有一些错误):

// returns an asyncIterator that will iterate over the observable values
function asyncIterator(observable) {
  const queue = []; // holds observed values not yet delivered
  let complete = false;
  let error = undefined;
  const promiseCallbacks = [];

  function sendNotification() {
    // see if caller is waiting on a result
    if (promiseCallbacks.length) {
      // send them the next value if it exists
      if (queue.length) {
        const value = queue.shift();
        promiseCallbacks.shift()[0]({ value, done: false });
      }
      // tell them the iteration is complete
      else if (complete) {
        while (promiseCallbacks.length) {
          promiseCallbacks.shift()[0]({ done: true });
        }
      }
      // send them an error
      else if (error) {
        while (promiseCallbacks.length) {
          promiseCallbacks.shift()[1](error);
        }
      }
    }
  }

  observable.subscribe(
    value => {
      queue.push(value);
      sendNotification();
    },
    err => {
      error = err;
      sendNotification();
    },
    () => {
      complete = true;
      sendNotification();
    });

  // return the iterator
  return {
    next() {
      return new Promise((resolve, reject) => {
        promiseCallbacks.push([resolve, reject]);
        sendNotification();
      });
    }
  }
}

与 for-wait-of 语言功能一起使用:

async someFunction() {
  const obs = ...;
  const asyncIterable = {
    [Symbol.asyncIterator]: () => asyncIterator(obs)
  };

  for await (const value of asyncIterable) {
    console.log(value);
  }
}

使用没有等待等待的语言功能:

async someFunction() {
  const obs = ...;
  const it = asyncIterator(obs);

  while (true) {
    const { value, done } = await it.next();
    if (done) {
      break;
    }

    console.log(value);
  }
}

【讨论】:

  • 不错!这很好用。您是否对此进行了编码,如果没有,您可以链接源代码吗?谢谢!
【解决方案2】:

这可能只是工作,因为 take(1) 完成了 observable,然后它被 await 消耗,下一个将产生对 value2 的第二个发射。

const observer= new Subject()
async function getStream(){
  const value1 = await observer.pipe(take(1)).toPromise() // 0.25
  const value2 = await observer.pipe(take(1)).toPromise() // 0.75
  return [value1,value2]
}
getStream().then(values=>{
  console.log(values)
})
//const obs = new Observable((observer) => {
 setTimeout(()=>{observer.next(0.25)},1000);
 setTimeout(()=>observer.next(0.75),2000);

更新:使用主题来发射。

【讨论】:

  • 这行不通。所有变量 (value1..value4) 将具有第一个发出的值 - 0.25。要使其正常工作,需要将每次拍摄更改为 (take(1)...take(4))。
  • 那么就在while循环中等待它们?或者只是订阅一个 observable,因为它们是要被使用的。或者根本不使用 observables,一开始就使用 Promise。
  • 更新了答案,已经测试过了,它会起作用的。虽然我认为它错过了异步编程的重点
  • @taras-d 是对的,您需要递增每个take,否则您将始终使用take(1) 获得第一个值。我使用Subject 测试了这里提出的解决方案,但这并不是一个好的选择,因为它不会等到客户端订阅后才推送值。可以订阅一个主题并且不会收到以前推送的值,因为该过程在客户端订阅之前已经启动。出于同样的原因,ReplaySubject 也不是一个好的选择,推送值的过程应该只在客户端订阅 observable 之后开始,而不是之前。
  • 那为什么不用生成器来代替
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2022-08-19
  • 2019-03-01
  • 2016-07-18
  • 2018-07-05
  • 2021-11-05
  • 1970-01-01
  • 2016-01-17
相关资源
最近更新 更多