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

如何创建函数检测Observable流中被过滤操作符排除的值?

检测Observable中被过滤操作符排除的值

嘿,这个需求我之前也碰到过——明明知道有些值被filter、debounceTime这类操作符拦下来了,但就是看不到它们,毕竟这些值只是被排除出输出流,并没有被转换或者销毁。下面给你一个基于RxJS的TypeScript实现方案,完美解决这个问题。

核心思路

我们的目标是:在应用过滤操作符的同时,把“通过的流”和“被排除的流”分开返回,这样你就能同时处理正常输出和被过滤掉的那些值。

通用实现方案

先给你一个通用的applyFilter函数,它能适配大多数过滤类操作符,包括同步的(比如filter)和异步的(比如debounceTime、throttleTime):

import { Observable, OperatorFunction, partition, map, share } from 'rxjs';

// 先定义一个带标记的类型,用来区分通过/排除的值
type FilteredResult<T> = 
  | { type: 'passed'; value: T }
  | { type: 'excluded'; value: T };

// 包装原过滤操作符,让它能标记每个值的状态
function wrapFilter<T>(operator: OperatorFunction<T, T>): OperatorFunction<T, FilteredResult<T>> {
  return (source) => new Observable(subscriber => {
    // 先共享源流,避免重复订阅
    const sharedSource = source.pipe(share());

    // 订阅通过过滤的流,标记为passed
    const passedSub = sharedSource.pipe(operator).subscribe({
      next: val => subscriber.next({ type: 'passed', value: val }),
      error: err => subscriber.error(err),
      complete: () => subscriber.complete()
    });

    // 订阅源流,找出那些没出现在passed流里的值,标记为excluded
    const seenPassed = new Set<T>();
    const excludedSub = sharedSource.subscribe({
      next: val => {
        // 加个微任务延迟,确保先处理passed流的订阅逻辑
        queueMicrotask(() => {
          if (!seenPassed.has(val)) {
            subscriber.next({ type: 'excluded', value: val });
          }
        });
      },
      error: err => subscriber.error(err)
    });

    // 记录通过的值
    const trackPassedSub = sharedSource.pipe(operator).subscribe(val => seenPassed.add(val));

    return () => {
      passedSub.unsubscribe();
      excludedSub.unsubscribe();
      trackPassedSub.unsubscribe();
    };
  });
}

// 最终的applyFilter函数
export function applyFilter<T>(
  source: Observable<T>,
  filterOperator: OperatorFunction<T, T>
): { passed: Observable<T>; excluded: Observable<T> } {
  const wrappedStream = source.pipe(wrapFilter(filterOperator));
  // 拆分出passed和excluded流
  const [passedStream, excludedStream] = partition(wrappedStream, res => res.type === 'passed');

  // 映射回原始值类型
  return {
    passed: passedStream.pipe(map(res => res.value)),
    excluded: excludedStream.pipe(map(res => res.value))
  };
}

使用示例

比如用debounceTime来测试,看看哪些值被防抖操作给忽略了:

import { interval } from 'rxjs';
import { debounceTime } from 'rxjs/operators';

// 每100ms发出一个递增数字
const source$ = interval(100);

// 应用300ms的防抖,同时获取被排除的流
const { passed, excluded } = applyFilter(source$, debounceTime(300));

// 订阅通过的流:只有当300ms内没有新值时,才会发出最后一个值
passed.subscribe(val => console.log(`✅ 通过的值: ${val}`));

// 订阅被排除的流:所有被防抖忽略的中间值
excluded.subscribe(val => console.log(`❌ 被排除的值: ${val}`));

注意事项

  • 对于引用类型的值,Set是基于引用相等性判断的,如果需要比较内容,你可以把seenPassed换成数组,用自定义的比较函数(比如JSON.stringify或者lodash的isEqual)来检查。
  • 对于像distinctUntilChanged这类依赖历史状态的操作符,上面的方案可能需要微调——因为这类操作符的判断依赖之前的值,这时候你可能需要直接在wrapFilter里重写操作符逻辑,同时跟踪通过和排除的状态。
  • 这个方案用了share()来避免重复订阅源流,性能上大部分场景都没问题,如果是超高频率的流,可以考虑进一步优化。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 08:18:24