RxJS中Observable发射值时重置audit操作符的实现问题
RxJS 实现c$发射时跳过audit节流直接执行计算
现有逻辑与需求
- 现有实现通过
combineLatest组合两个Observable:a$每秒发射4次值,基于当前b$的值配置audit操作符做节流控制 - 节流处理后通过
withLatestFrom获取c$的值,结合a$与c$的返回值完成业务计算 - 核心需求:当
c$发射值时直接跳过audit节流逻辑,无延迟进入计算环节,此前尝试的多种基于switchMap的实现均未达到预期
原有实现代码
combineLatest(a$, b$).pipe( audit((_, b) => { if (b == 1) { return interval(5000); } return interval(1000); }), withLatestFrom(c$) ) .subscribe(([[a, _], c]) => { // 基于a和c完成业务计算 });
可行实现方案
核心思路是拆分两个独立的触发通道做合并:保留原有节流触发的常规通道,新增c$触发的即时通道,两个通道输出统一格式的依赖值后直接进入计算环节。
import { merge, combineLatest, interval } from 'rxjs'; import { audit, map, switchMap, take } from 'rxjs/operators'; // 常规节流触发通道:完全保留原有audit节流规则 const throttledChannel$ = combineLatest(a$, b$).pipe( audit(([_, b]) => b === 1 ? interval(5000) : interval(1000)), switchMap(([a, b]) => c$.pipe(take(1), map(c => ({a, c}))) ) ); // c$即时触发通道:c$发射时直接取最新a、b值,完全绕过audit节流 const cTriggerChannel$ = c$.pipe( switchMap(c => combineLatest(a$, b$).pipe( take(1), map(([a]) => ({a, c})) ) ) ); // 合并两个通道,统一进入计算逻辑 merge(throttledChannel$, cTriggerChannel$) .subscribe(({a, c}) => { // 直接使用a、c完成业务计算 });
实现说明
- 两个触发通道完全独立,
c$触发时不会经过audit的延迟判断,拿到对应时刻最新的a值后立刻进入计算 - 常规节流通道逻辑和原有实现完全一致,不会破坏基于
b$配置的节流规则 - 两个通道输出的数据格式完全统一,计算逻辑不需要做分支判断
内容的提问来源于stack exchange,提问作者antonm76
相关产品推荐
相关产品推荐

