如何创建函数检测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
相关产品推荐
相关产品推荐

