You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RxJS技术需求:基于add*事件触发延迟/即时emit remove*事件

RxJS Solution for Add/Remove Event Pair Logic

I get exactly what you're trying to achieve here—each addN event needs a corresponding removeN emitted either after a delay (if no new add events come in) or immediately when a new add event arrives. Here's a clean implementation using RxJS operators, plus a complete marble test to validate the behavior.

Core Implementation

The key here is using switchMap to cancel the pending delay for the previous add event whenever a new one comes in. We'll create a custom inner observable that either emits the remove event after the delay, or emits it immediately if unsubscribed (which happens when switchMap switches to a new inner observable for the next add event).

import { Observable, switchMap } from 'rxjs';

// Define your delay duration (adjust as needed)
const DELAY = 1000; // 1 second

// Assume addEvents$ is your source observable emitting add events (e.g., { type: 'add', id: 1 })
const removeEvents$ = addEvents$.pipe(
  switchMap((addEvent) => {
    const removeEvent = { type: 'remove', id: addEvent.id };
    
    return new Observable((subscriber) => {
      // Schedule the delayed remove emission
      const timeoutId = setTimeout(() => {
        subscriber.next(removeEvent);
        subscriber.complete();
      }, DELAY);

      // Teardown logic: if we unsubscribe before the delay ends, emit remove immediately
      return () => {
        clearTimeout(timeoutId);
        subscriber.next(removeEvent);
        subscriber.complete();
      };
    });
  })
);

How It Works

  1. When an addN event is received, we schedule its corresponding removeN to be emitted after DELAY.
  2. If a new add* event arrives before the delay completes, switchMap unsubscribes from the previous inner observable. This triggers the teardown function, which clears the pending timeout and emits removeN immediately.
  3. switchMap then subscribes to a new inner observable for the latest add event, starting its own delay timer.
  4. If no new add events come in before the delay ends, the scheduled removeN is emitted as expected.

Complete Marble Test

Here's the finished test case using RxJS's TestScheduler to verify the behavior matches your requirements:

import test from 'tape';
import { TestScheduler } from 'rxjs/testing';
import { Observable, switchMap } from 'rxjs';

test('add/remove event pairing with delay and immediate cancellation', (t) => {
  const scheduler = new TestScheduler((actual, expected) => {
    t.deepEqual(actual, expected);
  });

  scheduler.run(({ cold, expectObservable }) => {
    const TEST_DELAY = 10; // Use 10 time units for marble testing
    const addEvents = cold('a 5 b 15 c', {
      a: { type: 'add', id: 1 },
      b: { type: 'add', id: 2 },
      c: { type: 'add', id: 3 }
    });

    const removeEvents = addEvents.pipe(
      switchMap((addEvent) => {
        const removeEvent = { type: 'remove', id: addEvent.id };
        
        return new Observable((subscriber) => {
          const timeoutId = setTimeout(() => {
            subscriber.next(removeEvent);
            subscriber.complete();
          }, TEST_DELAY);

          return () => {
            clearTimeout(timeoutId);
            subscriber.next(removeEvent);
            subscriber.complete();
          };
        });
      })
    );

    // Expected timeline:
    // - add1 at 0 → remove1 emitted immediately at 5 when add2 arrives
    // - add2 at 5 → remove2 emitted at 15 (delay completes, no new add before then)
    // - add3 at 20 → remove3 emitted at 30 (delay completes)
    expectObservable(removeEvents).toBe('5ms r1 10ms r2 15ms r3', {
      r1: { type: 'remove', id: 1 },
      r2: { type: 'remove', id: 2 },
      r3: { type: 'remove', id: 3 }
    });
  });

  t.end();
});

Test Scenario Breakdown

  • add1 is emitted at time 0: we schedule remove1 for time 10.
  • add2 arrives at time5: this cancels the pending remove1, which is emitted immediately at time5. We then schedule remove2 for time15.
  • No new add events arrive before time15: remove2 is emitted as scheduled.
  • add3 arrives at time20: we schedule remove3 for time30, which is emitted since no further add events come in.

You can extend this test with additional scenarios (like back-to-back add events, or a single add event with no follow-ups) to cover all edge cases.

内容的提问来源于stack exchange,提问作者neezer

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 09:47:33