【问题标题】:Can RxJS be used in a pull-based way?RxJS 可以以基于拉取的方式使用吗?
【发布时间】:2015-11-11 22:47:47
【问题描述】:

examples in the RxJS README 似乎暗示我们必须订阅来源。换句话说:我们等待源发送事件。从这个意义上说,来源似乎是基于推送的:来源决定何时创建新项目。

然而,这与迭代器形成鲜明对比,迭代器严格来说只有在请求时才需要创建新项目,即当调用next() 时。这是pull-based 行为,也称为 lazy 生成。

例如,一个流可以返回所有 Wikipedia 页面的素数。这些项目仅在您要求时才会生成,因为预先生成所有这些项目是一项相当大的投资,而且可能只读取其中的 2 或 3 个。

RxJS 是否也可以具有这种基于拉取的行为,以便仅在您请求时才生成新项目?

page on backpressure 似乎表明这还不可能。

【问题讨论】:

  • 我会让专家回答,但只是一个简短的说明。调用next 并不意味着新项目仅在您说的必要时创建,它只是意味着当时请求它们(并且作为迭代器同步提供)。典型的例子是基于数组的迭代器。在调用迭代器枚举其值之前,数组(以及其中的值)已经存在。实际上,您可以将数组视为保存值的缓冲区,而您的 next 运算符的作用就像带有 request(1) 的受控 Rx.Observable。
  • 完全正确——我稍微改写了措辞以反映这一点。但是,我不会说迭代器通常基于数组。更有趣的用例是迭代器不是基于数组的。

标签: stream reactive-programming rxjs


【解决方案1】:

简短的回答是否定的。

RxJS 是为 reactive 应用程序设计的,因此正如您已经提到的,如果您需要基于拉取的语义,您应该使用 Iterator 而不是 ObservableObservables设计为迭代器的基于推送的对应物,因此它们在算法上确实占据了不同的空间。

显然,我不能说这将永远发生,因为这是社区将决定的事情。但据我所知 1) 这种情况下的语义不是那么好 2) 这与 reacting 对数据的想法背道而驰。

可以在here 找到一个很好的概要。它适用于 Rx.Net,但这些概念同样适用于 RxJS。

【讨论】:

  • 谢谢——不过,这里有一个澄清:Iterators 仅适用于 同步 基于拉取的场景。尽管有comparison table,但我不会将Observable 称为Iterable 的确切异步对应物,因为迭代器是基于拉的,但RxJS 对象是基于推的。
【解决方案2】:

Controlled observable 从您引用的页面中可以将可观察到的推送更改为拉取。

var controlled = source.controlled();

// this callback will only be invoked after controlled.request()
controlled.subscribe(n => {
  log("controlled: " + n);
  // do some work, then signal for next value
  setTimeout(() => controlled.request(1), 2500);
});

controlled.request(1);

真正的同步迭代器是不可能的,因为它会在源不发射时阻塞。

在下面的sn-p中,受控订阅者在发出信号时只得到一个项目,并且不跳过任何值。

var output = document.getElementById("output");
var log = function(str) {
  output.value += "\n" + str;
  output.scrollTop = output.scrollHeight;
};

var source = Rx.Observable.timer(0, 1000);
source.subscribe(n => log("source: " + n));

var controlled = source.controlled();
// this callback will only be invoked after controlled.request()
controlled.subscribe(n => {
  log("controlled: " + n);
  // do some work, then signal for next value
  setTimeout(() => controlled.request(1), 2500);
});
controlled.request(1);
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/2.5.2/rx.all.js"></script>

<body>
  <textarea id="output" style="width:150px; height: 150px"></textarea>
</body>

【讨论】:

  • 但是在这种情况下,缓冲区正在源权中建立,即源仍然可能生成比消费者处理能力更多的项目?因为我正在寻找一种仅在请求时创建新项目的解决方案。
  • 冷的 observable (hot vs cold observables) 在被告知之前不会开始推送值。如果这不符合您的要求,示例代码将有助于澄清问题。
  • 如果你想直接控制元素的生成,那么 observable 不是你想要的。观察者模式的全部要点是订阅者(观察者)不控制发布者的行为,他们只听。有多种方法可以实现异步迭代器,ES6 生成器是最明显的。像 Babeljs 这样的编译器现在可以让你在代码中使用生成器。
  • 您似乎将迭代器与生成器混淆了。 ES6 生成器是迭代器的异步版本。对于生成器function *gen() { yield 'a'; yield 'b'; },不会计算第二行,除非您在gen() 返回的生成器上调用next() 两次。在第一个yield 之后,控制权立即交还给调用者。
  • 这不是异步的;这是通过更早地返回控制来暂停执行。当元素需要时间来生成并且您不想主动等待它们时,异步是必要的。例如,以一个迭代器为例,它返回维基百科中每个页面的内容。迭代器需要执行一个 JSON 请求,交还控制权,并且只有当内容已经下载后,才返回页面。如果您连续调用下 100 次,则不能指望已经收到 100 页,因为它们还没有准备好。这就是我想做的事情。 (同时写了我自己的库。)
【解决方案3】:

我来晚了,但实际上将生成器与可观察对象结合起来非常简单。您可以通过将其与源 observable 同步来从生成器函数中提取值:

const fib = fibonacci()
interval(500).pipe(
  map(() => fib.next())
)
  .subscribe(console.log)

生成器实现供参考:

function* fibonacci() {
  let v1 = 1
  let v2 = 1
  while (true) {
    const res = v1
    v1 = v2
    v2 = v1 + res
    yield res
  }
}

【讨论】:

    猜你喜欢
    • 2021-09-17
    • 1970-01-01
    • 1970-01-01
    • 2013-07-12
    • 2018-07-28
    • 1970-01-01
    • 2019-12-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多