【问题标题】:Delaying all items except specific one延迟除特定项目外的所有项目
【发布时间】:2020-10-14 19:15:49
【问题描述】:

假设我有一个动作流。它们是提示、响应(对提示)或效果。它们以不规则的间隔出现,但假设每个间隔有 1 秒的延迟。

在每个PROMPT 操作上,我想发出该操作和一个 BEGIN 操作(假设我们想向用户显示消息 N 秒)。所有其他项目应延迟 N 秒,然后触发 END 操作(隐藏消息),一切继续。

这是我的代码(https://rxviz.com/):

const { interval, from, zip, timer } = Rx;
const { concatMap, delayWhen } = RxOperators;

const PROMPT = 'P';
const RESPONSE = 'R';
const EFFECT = 'E';

const BEGIN = '^';
const END = '&';

const convertAction = action => (action === PROMPT) ? [PROMPT, BEGIN, END] : [action];

// Just actions coming at regular intervals
const action$ = zip(
  from([PROMPT, RESPONSE, EFFECT, PROMPT, RESPONSE, EFFECT, EFFECT, EFFECT]),
    interval(1000),
  (a, b) => a,
);

action$.pipe(
  concatMap(action =>
    from(convertAction(action)).pipe(
      delayWhen(action => (action == END) ? timer(5000) : timer(0)),
    ),
  ),
);

我真正想做的是在PROMPT 之后的第一个RESPONSE 操作不受延迟的影响。如果它出现在 END 动作之前,它应该立即显示。所以,而不是

P^ &REP^ &REEE

我想收到

P^ R &EP^R &EEE

如何在将每个RESPONSE 保留在其对应的PROMPT 之后实现它?假设在PROMPTRESPONSE 之间没有任何事件。

【问题讨论】:

  • 你在zip一个函数? (a, b) => a - 我很惊讶帽子的作用。那有什么作用?
  • 它采用数组的值并将它们按间隔时间间隔。 a 是第一个流项目(元素),b 是第二个流项目(延迟无值)。
  • TIL 如果有一个函数作为最后一个参数,zip 使用它来格式化它的输出......整洁! -- 我一直用zip().pipe(map()) 来达到同样的效果。
  • 是否可以假设第二个提示总是在第一个响应之后,或者更一般地说,Promp N+1 总是在 ResponseN?
  • @Picci 是的。它们基本上无法区分,但每个 Prompt 后面都会有 Response(在未知时间之后)。

标签: javascript rxjs reactive-programming


【解决方案1】:

如果我理解正确的话,这是一个用 Observables 流来解决的非常有趣的问题。这就是我要攻击它的方式。

首先,我将在每个 PROMPTBEGINEND 操作除以延迟之后,将原始逻辑的结果存储在常量 actionDelayed$ 中,即我们引入的流。

const actionDelayed$ = action$.pipe(
  concatMap(action =>
    from(convertAction(action)).pipe(
      delayWhen(action => (action == END) ? timer(5000) : timer(0)),
    ),
  ),
);

然后我将创建 2 个单独的流,response$promptDelayed$,仅包含引入延迟之前的 RESPONSE 操作和引入延迟之后的 PROMPT 操作,像这样

const response$ = action$.pipe(
  filter(a => a == RESPONSE)
)
const promptDelayed$ = actionDelayed$.pipe(
  filter(a => a == PROMPT)
)

有了这 2 个流,我可以在发出 PROMPT 延迟动作之后创建另一个 RESPONSE 动作流,就像这样

const responseN1AfterPromptN$ = zip(response$, promptDelayed$).pipe(
  map(([r, p]) => r)
)

此时我只需像这样从actionDelayed$ 中删除所有RESPONSE 操作

const actionNoResponseDelayed$ = actionDelayed$.pipe(
  filter(a => a != RESPONSE)
)

并将actionNoResponseDelayed$responseN1AfterPromptN$ 合并以获得最终流。

要使用rxviz 尝试的全部代码是这样的

const { interval, from, zip, timer, merge } = Rx;
const { concatMap, delayWhen, share, filter, map } = RxOperators;

const PROMPT = 'P';
const RESPONSE = 'R';
const EFFECT = 'E';

const BEGIN = '^';
const END = '&';

const convertAction = action => (action === PROMPT) ? [PROMPT, BEGIN, END] : [action];

// Just actions coming at regular intervals
const action$ = zip(
  from([PROMPT, RESPONSE, EFFECT, PROMPT, RESPONSE, EFFECT, EFFECT, EFFECT]),
    interval(1000),
  (a, b) => a,
).pipe(share());

const actionDelayed$ = action$.pipe(
  concatMap(action =>
    from(convertAction(action)).pipe(
      delayWhen(action => (action == END) ? timer(5000) : timer(0)),
    ),
  ),
  share()
);

const response$ = action$.pipe(
  filter(a => a == RESPONSE)
)
const promptDelayed$ = actionDelayed$.pipe(
  filter(a => a == PROMPT)
)
const responseN1AfterPromptN$ = zip(response$, promptDelayed$).pipe(
  map(([r, p]) => r)
)
const actionNoResponseDelayed$ = actionDelayed$.pipe(
  filter(a => a != RESPONSE)
)

merge(actionNoResponseDelayed$, responseN1AfterPromptN$)

在创建action$actionDelayed$ 流时使用share 运算符可以避免在创建解决方案中使用的后续流时重复订阅这些流。

【讨论】:

    【解决方案2】:

    由于您使用的是concatMap,因此可能无法以这种方式工作。如您所知,它会等待complete 的内部可观察对象,然后再开始处理(订阅)待处理的对象。它在内部使用一个缓冲区,这样如果内部 observable 仍然处于活动状态(不是complete),则发出的值将被添加到该缓冲区。当内部 observable 变为非活动状态时,将选择缓冲区中最旧的值,并根据提供的回调函数创建 new 内部 observable。

    还有delayWhen,它会在其所有待处理 observables 完成后发出完整通知:

    // called when an inner observable sends a `next`/`complete` notification
    const notify = () => {
      // Notify the consumer.
      subscriber.next(value);
    
      // Ensure our inner subscription is cleaned up
      // as soon as possible. Once the first `next` fires,
      // we have no more use for this subscription.
      durationSubscriber?.unsubscribe();
    
      if (!closed) {
        active--;
        closed = true;
        checkComplete();
      }
    };
    

    checkComplete() 将检查是否需要向主流发送complete 通知:

    const checkComplete = () => isComplete && !active && subscriber.complete();
    

    我们发现activenotify() 中减少。当主源完成时,isComplete 变为 true

    // this is the `complete` callback
    () => {
      isComplete = true;
      checkComplete();
    }
    

    所以,这就是它不能以这种方式工作的原因:

    • PROMPT 操作用于创建concatMap 的第一个内部可观察对象
    • observable 发出 3 个连续动作 [PROMPT, BEGIN, END]
    • 前 2 个将获得 timer(0),而第三个 END 将获得 (timer(5000));请注意,此时,在发出 PROMPT 操作之前,isComplete 变量设置为 true,因为在这种情况下 from() 同步完成
    • 所以有一个timer(5000) 保留内部obs。积极的;然后从actions$ 流中发出RESPONSE,但由于还没有位置,它将被添加到缓冲区和内部obs。将在timer(5000) 最终过期时创建

    解决此问题的一种方法可能是将concatMap 替换为mergeMap

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-11-07
      • 1970-01-01
      • 2020-05-04
      • 2016-12-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多