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

RxJS问题:多次订阅后distinct操作符失效

Understanding RxJS's distinct Operator Behavior Across Subscriptions

Hey there, let's break down why you're seeing this behavior with distinct—it's a common gotcha tied to how RxJS handles subscriptions and cold vs. hot observables.

The Core Issue: distinct Resets Its State Per Subscription

Here's the key thing to remember: every time you subscribe to an observable that uses distinct, the operator creates a new internal set to track unique values. This set is tied directly to that specific subscription, not the observable itself.

Let's use a simple example that matches your scenario:

import { of } from 'rxjs';
import { distinct } from 'rxjs/operators';

// This is a COLD observable—every subscription re-emits all values from scratch
const coldSource$ = of(1, 2, 2, 3);

// First subscription: outputs 1, 2, 3 (duplicate 2 is filtered out)
coldSource$.pipe(distinct()).subscribe(val => console.log('Sub 1:', val));

// Second subscription: outputs 1, 2, 3 AGAIN
// Why? Because this subscription gets a brand new instance of `distinct` with an empty set
coldSource$.pipe(distinct()).subscribe(val => console.log('Sub 2:', val));

Cold observables (like of(), from(), or HTTP requests) re-run their entire emission sequence for every new subscription. So each time you subscribe, distinct starts fresh with no memory of previous values from other subscriptions.

How to Fix It: Share the distinct State Across Subscriptions

If you want distinct to maintain its unique value tracking across multiple subscriptions, you need to make sure all subscriptions share the same instance of the distinct operator. There are two main ways to do this:

1. Use a Hot Observable with share()

Convert your cold observable to a hot one using operators like share(), which multicasts the emission to all subscribers. This way, the distinct operator is only initialized once, and its internal set is shared:

import { of } from 'rxjs';
import { distinct, share } from 'rxjs/operators';

// Now a HOT observable—subscriptions share the same emission stream
const hotSource$ = of(1, 2, 2, 3).pipe(share());

// First subscription: outputs 1, 2, 3
hotSource$.pipe(distinct()).subscribe(val => console.log('Sub 1:', val));

// Second subscription: If the source has already completed, it won't emit anything
// For an ongoing stream (like interval), new subscribers will only get values that haven't been seen yet

For a continuous stream example:

import { interval } from 'rxjs';
import { distinct, share, map } from 'rxjs/operators';

// Emits 0,1,2,0,1,2,... every second
const ongoingStream$ = interval(1000).pipe(
  map(i => i % 3),
  share() // Multicast to all subscribers
);

const distinctStream$ = ongoingStream$.pipe(distinct());

// First subscriber gets 0,1,2, then nothing else (all duplicates after that)
distinctStream$.subscribe(val => console.log('Sub 1:', val));

// Subscribe 1.5 seconds later—this subscriber will get 2, then nothing else
setTimeout(() => {
  distinctStream$.subscribe(val => console.log('Sub 2:', val));
}, 1500);

2. Global Permanent Deduplication

If you need to track unique values across all subscriptions, even after the stream completes and re-runs, you'll need to maintain an external set and use filter instead:

import { Observable } from 'rxjs';
import { filter } from 'rxjs/operators';

// Global set to track all unique values forever
const globalUniqueValues = new Set<number>();

// A stream that emits incrementing values each time it's subscribed to
const dynamicSource$ = new Observable(subscriber => {
  let nextVal = globalUniqueValues.size + 1;
  subscriber.next(nextVal++);
  subscriber.next(nextVal++);
  subscriber.next(nextVal++);
  subscriber.complete();
});

// Filter out values already in the global set
const globallyDistinct$ = dynamicSource$.pipe(
  filter(val => {
    if (!globalUniqueValues.has(val)) {
      globalUniqueValues.add(val);
      return true;
    }
    return false;
  })
);

// First subscription: outputs 1,2,3
globallyDistinct$.subscribe(val => console.log('Sub 1:', val));

// Second subscription: outputs 4,5,6
globallyDistinct$.subscribe(val => console.log('Sub 2:', val));

Recap

  • distinct is subscription-scoped: it resets its tracking set for every new subscription to a cold observable.
  • To share deduplication state across subscriptions, use a hot observable with share() to multicast the stream.
  • For permanent global deduplication, manage your own external set and use filter.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:47:36