【问题标题】:Passing pipeline item to promise argument in `takeUntil`将管道项目传递给“takeUntil”中的承诺参数
【发布时间】:2021-08-11 14:12:30
【问题描述】:

我有与这个例子类似的控制流的代码(显然,下面的谓词不需要是async,但它是一个例子):

const items [1,2,3,4,5];
const predicate = async (i) => i < 3;
const pipeline = from(items).pipe(
  takeUntil(predicate),
);

pipeline.subscribe(console.log);

但这会引发 TypeError 消息“您可以提供 Observable、Promise、ReadableStream、Array、AsyncIterable 或 Iterable。”

我尝试过让predicate 成为一个承诺(new Promise(...),并使用takeWhile 代替takeUntil,但都没有按预期工作(承诺总是返回真实 - 我假设它被强制为真实). 这是我对 takeUntil/takeWhile 工作方式的某种误解吗?

作为一种解决方法,我目前正在使用这个:

const takeWhileAsync = (predicate = tautology) => {
  const resultSymbol = Symbol('predicateResult');
  const valueSymbol = Symbol('value');

  const predicateResolver = item => of(item).pipe(
    concatMap(async (i) => {
      const predicateResult = await predicate(i);
      return {[resultSymbol]: predicateResult, [valueSymbol]: i};
    }),
  );

  return pipe(
    concatMap(predicateResolver),
    takeWhile(({[resultSymbol]: predicateResult}) => predicateResult),
    pluck(valueSymbol),
  );
};

【问题讨论】:

  • takeUntil 不接受每个发出的项目都要调用的谓词,它接受一个“通知器”,一个在停止时发出的可观察对象。
  • @jonrsharpe 有没有一种方法可以提前从管道中退出,而无需创建可观察的通知程序?理想情况下,我想传入一个谓词(与takeWhile 一样),但该谓词是一个异步函数(takeWhile 似乎总是强制为真值)
  • 一个异步函数应该更正为一个返回承诺的函数,对不起
  • 据我所知不是这样;正如您所看到的,takeWhile 需要一个返回布尔值的谓词,而不是一个诺言(这确实是真的)。​​

标签: javascript rxjs


【解决方案1】:

惯用的 RxJS

大多数 RxJS 操作符(concatMapmergeMapswitchMap 等...)将 ObservableInput 作为返回值(这意味着它们可以原生使用 Promises)。

这是对@'Nick Bull 的回答的一个看法,它没有做任何promise (async/await) 的事情。通过这种方式,您可以使用 Promise 或(可能是可取的)完全坚持使用 Observables。

function takeWhileConcat<T>(genPred: (v:T) => ObservableInput<Boolean>): MonoTypeOperatorFunction<T>{
  return pipe(
    concatMap((payload: T) => from(genPred(payload)).pipe(
      take(1),
      map((pass: boolean) => ({payload, pass}))
    )),
    takeWhile(({pass}) => pass),
    map(({payload}) => payload)
  );
}

const items = [1,2,3,4,5];
const predicate = async (i) => i < 3;
const pipeline = from(items).pipe(
  takeWhileConcat(predicate),
);

pipeline.subscribe(console.log);

现在,如果您愿意,可以将谓词替换为可观察对象:

const predicate = i =&gt; of(i &lt; 3);

不要改变任何其他东西。这很好,因为 observables 和 Promise 有时并没有你想象的那么好。

考虑到 Promise 是急切的,而 observables 是惰性的,你可能会得到一些难以调试的奇怪执行顺序。


此解决方案不允许使用非承诺谓词!

所以,你是对的。此解决方案要求您返回 ObservableInput(任何可迭代、承诺或可观察的)。实际上,任何 ES6 可迭代对象,如数组、生成器、映射、哈希映射、向量、自定义可迭代对象,应有尽有。它们都会起作用。

  • 可观察:predicate = value =&gt; of(value &gt; 3)
  • 可迭代:predicate = value =&gt; [value &gt; 3]
  • 承诺:predicate = value =&gt; Promise.resolve(value &gt; 3)
  • 承诺的语法糖:predicate = async value =&gt; value &gt; 3

不允许允许的是不是ObservableInput 的任何东西。这与其他所有采用 ObservableInput 函数的 RxJS 运算符的方式相匹配。当然,我们可以使用 of 将任何值作为 observable 进行军事化,但这被决定反对,因为它更有可能是一把枪而不是有用。

在动态类型语言中,很难确定您的 API 允许什么以及应该在哪里引发错误。我喜欢 RxJS 默认不将值视为 Observables。我认为 RxJS api 更清晰。

运营商在表达意图方面做得更好。想象一下这两个是一样的:

map(x =&gt; x + 1)
mergeMap(x = x + 1)

第二个可以将返回的值转换为可观察的并合并该可观察的,但这需要有关此运算符的大量专业知识。另一方面,Map 的工作方式与我们已经熟悉的其他迭代器/集合完全相同。

如何接受非承诺谓词

无论如何,如果您愿意,您可以更改我的答案以接受标准谓词 (v =&gt; boolean) 以及异步谓词 (v =&gt; ObservableInput&lt;boolean&gt;)。只需提供一个值并检查返回的内容。

我只是不相信这是可取的行为。

如果输入项是无限生成器怎么办?

这是一个永远生成整数的生成器。

const range = function*() { 
  for (let i = 0; true; i++) yield i; 
}

from(range()) 不知道何时停止调用生成器(甚至不知道生成器是无限的)。 from(range()).subscribe(console.log) 将无限期地将数字打印到控制台。

这里的关键是,在这种情况下,阻止我们回调生成器的代码必须同步运行。

例如:

from(range()).pipe(
  take(5)
).subscribe(console.log);

会将数字 0 - 4 打印到控制台。

对于我们的自定义运算符也是如此。仍然有效的代码:

from(range()).pipe(
  takeWhileConcat(v => of(v < 10))
).subscribe(console.log);

// or 

from(range()).pipe(
  takeWhileConcat(v => [v < 10])
).subscribe(console.log);

不会停止的代码:


from(range()).pipe(
  takeWhileConcat(v => of(v < 10).pipe(
    delay(0)
  ))
).subscribe(console.log);

// or

from(range()).pipe(
  takeWhileConcat(async v => v < 10)
).subscribe(console.log);

这是 javascript 引擎如何处理异步行为的结果。在引擎查看事件队列之前,任何当前代码都会运行完成。每个 Promise 都放在事件队列中,异步 observables 也放在事件队列中(这就是为什么 delay(0) 与立即解决的 Promise 基本相同)

concatMap 确实有内置的背压,但是代码的异步部分永远不会运行,因为代码的同步部分已经创建了一个无限循环。

这是基于推送的流媒体库(如 RxJS)的缺点之一。如果它是基于拉的(就像生成器一样),这不会是一个问题,但会出现其他问题。您可以在 Google 上搜索基于拉/推的流式传输以获取有关该主题的大量文章。

有一些安全的方法可以连接基于拉的流和基于推的流,但需要一些工作。

【讨论】:

  • 感谢您解释为什么不在这些类型的运算符中使用 Promise,这对新手来说真是个好知识!我知道坚持使用 observables 可能会更好,我将研究如何为我正在使用的一些基于 promise 的库编写适配器,以及如何将它们转换为 observables
  • 此解决方案不允许使用非承诺谓词!因此,需要用util.promisify(或任何其他承诺)包装任何非承诺。有没有办法解决这个问题,而不需要外部实用程序?
  • 如果您想使用非承诺谓词,只需使用takeWhile
  • 我更新了我的答案。它没有实现这一点,但有一个关于如何完成的提示。
  • @NickBull 我添加了另一个更新来讨论为什么会发生这种情况。
【解决方案2】:

根据您的评论:

理想情况下,我想传入一个谓词(与 takeWhile 一样),但该谓词是一个异步函数(takeWhile 似乎总是强制为真值)

我相信你可以这样做:

const items [1,2,3,4,5];
const predicate = async (i) => i < 3;
const pipeline = from(items).pipe(
  concatMap(predicate),
  filter(Boolean),
  take(1)
);

pipeline.subscribe(console.log);

通过这样做,您将链接所有异步调用 (concatMap),但可以随意在 // (mergeMap) 中运行它们,或者在有新值进入 (switchMap) 时取消它们。然后您应用过滤器同步,它只会让通过谓词的值和take(1) 将在通过谓词的第一个值之后结束流。

编辑:

在此评论之后,这是新的解决方案:

const items = [1, 2, 3, 4, 5];
const predicate = async i => i < 3;
const pipeline = from(items).pipe(
  concatMap(item =>
    from(predicate(item)).pipe(
      filter(Boolean),
      mapTo(item)
    )
  )
);

现场示例:https://stackblitz.com/edit/rxjs-yenhpy?devtoolsheight=60

【讨论】:

  • 这似乎正确地利用了谓词,但在快速测试后,输出是布尔结果,而不是值。我也不确定我是否需要take(1),因为我只想在谓词为假时退出。我是 RxJS 的新手,如果我误解了,请原谅我!
  • 只是为了进一步阐明谓词async (i) =&gt; i &lt; 3items = [1,2,3,4,5] 的要求,我想收到1 和2
  • 哦,我明白了,等等
  • @NickBull 我已经进行了编辑 :) 让我知道这是否是你想要的
  • 是的,按预期工作 - 感谢 Maxime!
【解决方案3】:

根据问题cmets,没有直接的方法可以使用takeWhile

这是上述问题的改进解决方法,没有多余的of.pipe(concatMap)

const takeWhileAsync = (predicate = tautology) => {
  const resultSymbol = Symbol('predicateResult');
  const valueSymbol = Symbol('value');

  const predicateResolver = item => async (i) => {
    const predicateResult = await predicate(i);
    return {[resultSymbol]: predicateResult, [valueSymbol]: i};
  };

  return pipe(
    concatMap(predicateResolver),
    takeWhile(({[resultSymbol]: predicateResult}) => predicateResult),
    pluck(valueSymbol),
  );
};

【讨论】:

    猜你喜欢
    • 2017-12-03
    • 2015-06-06
    • 2021-08-23
    • 2014-09-14
    • 1970-01-01
    • 2021-09-25
    • 2017-03-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多