RxJS技术需求:基于add*事件触发延迟/即时emit remove*事件
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
- When an
addNevent is received, we schedule its correspondingremoveNto be emitted afterDELAY. - If a new
add*event arrives before the delay completes,switchMapunsubscribes from the previous inner observable. This triggers the teardown function, which clears the pending timeout and emitsremoveNimmediately. switchMapthen subscribes to a new inner observable for the latest add event, starting its own delay timer.- If no new add events come in before the delay ends, the scheduled
removeNis 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

