【问题标题】:Filtering RxJS stream after emission of other observable until timer runs out在其他可观察的发射后过滤 RxJS 流,直到计时器用完
【发布时间】:2020-10-04 11:10:24
【问题描述】:

我想在 RxJS 中实现以下行为,但找不到使用可用运算符的方法:

  • 流 A:由连续的事件流生成(例如浏览器滚动)
  • 流 B:由另一个任意事件(例如某种用户输入)生成
  • B 发出一个值时,我想暂停 A 的处理,直到经过指定的时间。 A 在此时间范围内发出的所有值都将被丢弃。
  • B 在此间隔内发出另一个值时,将重置间隔。
  • 间隔过后,A 发出的值不再被过滤。
// Example usage.
streamA$
  .pipe(
    unknownOperator(streamB$, 800),
    tap(val => doSomething(val))
  )
// Output: E.g. [event1, event2, <skips processing because streamB$ emitted>, event10, ...]

// Operator API.
const unknownOperator = (pauseProcessingWhenEmits: Observable<any>, pauseIntervalInMs: number) => ...

我认为throttle 可以用于这个用例,但是它不会让任何发射通过,直到 B 第一次发射(可能永远不会!)。

streamA$
  .pipe(
    // If B does not emit, this never lets any emission of A pass through!
    throttle(() => streamB$.pipe(delay(800)), {leading: false}),
    tap(val => doSomething(val))
  )

一个简单的技巧是例如手动订阅 B,在 Angular 组件中存储发出值时的时间戳,然后过滤直到指定的时间过去:
(显然违背了反应式框架的副作用避免)

streamB$
  .pipe(
    tap(() => this.timestamp = Date.now())
  ).subscribe()

streamA$
  .pipe(
    filter(() => Date.now() - this.timestamp > 800),
    tap(val => doSomething(val))
  )

在我构建自己的自定义运算符之前,我想在这里与专家核实是否有人知道一个运算符(组合)可以在不引入副作用的情况下执行此操作:)

【问题讨论】:

    标签: javascript angular rxjs


    【解决方案1】:

    我认为这是一种方法:

    bModified$ = b$.pipe(
      switchMap(
        () => of(null).pipe(
          delay(ms),
          switchMapTo(subject),
          ignoreElements(),
          startWith(null).
        )
      )
    )
    
    a$.pipe(
      multicast(
        new Subject(),
        subject => merge(
          subject.pipe(
            takeUntil(bModified$)
          ),
          NEVER,
        )
      ),
      refCount(),
    )
    

    这似乎不是一个问题,其解决方案必然涉及多播,但在上述方法中,我使用了一种本地多播

    这不是预期的多播行为,因为如果您订阅a$ 多次(假设N 次),源将到达N 次,所以多播确实不会在该级别发生

    所以,让我们检查每个相关部分:

    multicast(
      new Subject(),
      subject => merge(
        subject.pipe(
          takeUntil(bModified$)
        ),
        NEVER,
      )
    ),
    

    第一个参数将指示为实现本地多播而使用的主题类型。第二个参数是一个函数,更准确地称为选择器。它的单个参数是之前指定的参数(Subject 实例)。每次订阅 a$ 时都会调用此选择器函数。

    我们可以从source code看到:

    selector(subject).subscribe(subscriber).add(source.subscribe(subject));
    

    源已订阅,source.subscribe(subject)。通过selector(subject).subscribe(subscriber) 实现的是一个新的subscriber,它将成为Subject 的观察者列表的一部分(它始终是相同的Subject 实例),因为merge 在内部订阅了提供的可观察对象。

    我们使用merge(..., NEVER) 是因为,如果订阅选择器的订阅者完成,那么下次a$ 流再次变为活动状态时,必须重新订阅源。通过附加 NEVER,调用 select(subject) 的 observable 结果表单将永远完成,因为为了让 merge 完成,它的所有 observables 都必须完成。

    subscribe(subscriber).add(source.subscribe(subject))subscribedSubject 之间创建连接,这样当subscriber 完成时,Subject 实例将调用其unsubscribe method

    所以,假设我们订阅了a$a$.pipe(...).subscribe(mySubscriber)。正在使用的Subject 实例将有一个订阅者,如果a$ 发出一些东西,mySubscriber 将接收它(通过主题)。

    现在让我们讨论bModified$ 发出的情况

    bModified$ = b$.pipe(
      switchMap(
        () => of(null).pipe(
          delay(ms),
          switchMapTo(subject),
          ignoreElements(),
          startWith(null).
        )
      )
    )
    

    首先,我们使用switchMap,因为一个要求是当b$ 发出时,计时器应该重置。但是,我看到这个问题的方式是,当b$ 发出时,两件事必须发生:

    • 启动计时器 (1)
    • 暂停a$的排放 (2)

    (1) 是通过在Subject 的订阅者中使用takeUntil 来实现的。通过使用startWithb$ 将立即发射,因此a$ 的发射被忽略。在switchMap 的内部可观察对象中,我们使用delay(ms) 来指定计时器应该花费多长时间。在它过去之后,在switchMapTo(subject) 的帮助下,Subject 现在将获得一个新的订阅者,这意味着a$ 的排放将被mySubscriber 接收(无需重新订阅源)。最后,使用ignoreElements,因为否则当a$发出时,这意味着b$也发出,这将导致a$再次停止。在switchMapTo(subject) 之后的是a$ 的通知。

    基本上,我们可以通过这种方式实现可暂停行为:当Subject 实例作为一个订阅者时(在此解决方案中它最多有一个),它没有暂停。如果没有,则表示它已暂停

    编辑:或者,您可以查看来自rxjs-etcpause operator

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-03-29
      • 1970-01-01
      • 2020-09-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多