【问题标题】:How can I create a pausableBuffer w/ rxjs 5如何使用 rxjs 5 创建一个 pausableBuffer
【发布时间】:2016-11-11 23:20:38
【问题描述】:

我正在尝试制作我认为的 pausable buffer

我有人为此分享了他们的代码,但我不知道如何将其转换为自定义操作(没有打字稿/只有 ES6。

const attach = Rx.Observable.timer(0 * 1000, 8 * 1000).mapTo('@');
const detach = Rx.Observable.timer(4 * 1000, 8 * 1000).mapTo('#');

const input = Rx.Observable.interval(1* 1000);
const pauser = attach.mapTo(true).merge(detach.mapTo(false));

input
  .publish(_input => _input
    .combineLatest(pauser, (v, b) => b)
    .filter(e => e)
    .publish(_switch => _input.bufferWhen(() => _switch.take(1)))
  )
  .flatMap(e => Rx.Observable.from(e))
  .concatMap(e => Rx.Observable.empty().delay(150).startWith(e))

有人可以帮我创建一个这样我就可以做到input.pausableBuffer(pauser)(也许定义一个startsWith)。

【问题讨论】:

标签: rxjs rxjs5


【解决方案1】:

你可以像这样将它添加到原型中:

var pausableBuffer = function(pauser) {
  return this.publish(_input => _input
    .combineLatest(pauser, (v, b) => b)
    .filter(e => e)
    .publish(_switch => _input.bufferWhen(() => _switch.take(1)))
  )
  .flatMap(e => Rx.Observable.from(e));
}

Rx.Observable.prototype.pausableBuffer = pausableBuffer;

要记住的一点是,这将从暂停状态开始。要改为在活动状态下启动它,请将.startWith(true) 添加到pauser

var pausableBuffer = function(pauser) {
  return this.publish(_input => _input
    .combineLatest(pauser.startWith(true), (v, b) => b)
    .filter(e => e)
    .publish(_switch => _input.bufferWhen(() => _switch.take(1)))
  )
  .flatMap(e => Rx.Observable.from(e));
}

Rx.Observable.prototype.pausableBuffer = pausableBuffer;

2019 年更新:RxJs 6 样式:

var pausableBuffer = function(pauser) {
  return (source) => source.pipe(publish(_input => 
  combineLatest(_input, pauser.pipe(startWith(true))).pipe(
    map(([inp, pa]) => pa),
    filter(pa => pa),
    publish(_switch => _input.pipe(bufferWhen(() => _switch.pipe(take(1)))))
  )),
    mergeMap(e => from(e))
  );
}

Demo

【讨论】:

  • 这个例子在 RxJS 6 中是什么样子的?
  • @chrismarx 实际上更简单的是,您添加管道,而不是将函数添加到原型中,而是使其成为返回函数的内部函数。内部函数的参数将替换示例中的 thispausableBuffer = (pauser) => (source) => source.publish(_input => _input.pipe(...))
  • 感谢您提供示例。我在使用这个答案时遇到了问题,因为 this.publish 不再有效,只是切换到“发布”也不起作用。我最终使用了这里的缓冲示例,它使用了 bufferToggle 和 windowToggle - medium.com/@kddsky/pauseable-observables-in-rxjs-58ce2b8c7dfd
  • @chrismarx publish 仍然是有效的运算符。我只是意识到我忘了pipe它。
  • @chrismarx 我在答案中添加了完整的 RxJs 6 解决方案 + JsBin 演示。
猜你喜欢
  • 2016-11-27
  • 2016-09-06
  • 1970-01-01
  • 1970-01-01
  • 2016-04-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-06-11
相关资源
最近更新 更多