RxJS技术问询:移除compose、转为BehaviorSubject及追踪触发流
针对你的RxJS问题的分步解决方案
我来逐个解决你提出的四个技术需求,直接上代码和易懂的解释:
问题1:移除Ramda的compose函数,仅用RxJS实现对象处理
原来的Ramda compose 是从右到左执行函数(addMiddleName → addLastName → addFirstName),我们完全可以用原生的函数链式调用或者对象扩展来替代,不需要依赖第三方库:
方案A:函数链式调用(和compose逻辑完全一致)
combineLatest([f1$, f2$, f3$, mainObj$], (f, l, m, dataBlock) => { // 等价于 R.compose(addFirstName(f), addLastName(l), addMiddleName(m))(dataBlock) return addFirstName(f)(addLastName(l)(addMiddleName(m)(dataBlock))); })
方案B:直接对象扩展(更直观,前提是addXxx函数是纯对象更新)
如果你的addFirstName这类函数只是给对象添加/更新属性,比如:
const addFirstName = (firstName: string) => (obj: any) => ({...obj, firstName});
那可以直接合并对象,代码更简洁:
combineLatest([f1$, f2$, f3$, mainObj$], (f, l, m, dataBlock) => { return { ...dataBlock, firstName: f, lastName: l, middleName: m }; })
问题2:让f1$、f2$、f3$表现得像BehaviorSubject,支持仅地址保存
BehaviorSubject的核心特性是有初始值,订阅时会立即发出当前最新值。要让你的事件流具备这个特性,有两种简单方式:
方案A:直接改用BehaviorSubject
// 初始化时给个默认值(比如null/空字符串,根据业务需求定) const f1$ = new BehaviorSubject<string | null>(null); const f2$ = new BehaviorSubject<string | null>(null); const f3$ = new BehaviorSubject<string | null>(null); // 绑定事件触发next this.eventService.on('firstName', data => f1$.next(data)); this.eventService.on('lastName', data => f2$.next(data)); this.eventService.on('middleName', data => f3$.next(data));
方案B:给原有Observable添加startWith操作符
如果不想修改事件绑定的逻辑,直接在原流后加startWith注入初始值即可:
import { startWith } from 'rxjs/operators'; const f1$ = Observable.create(ob => this.eventService.on('firstName', data => ob.next(data))) .pipe(startWith(null)); // 初始值根据业务调整 const f2$ = Observable.create(ob => this.eventService.on('lastName', data => ob.next(data))) .pipe(startWith(null)); const f3$ = Observable.create(ob => this.eventService.on('middleName', data => ob.next(data))) .pipe(startWith(null));
这样combineLatest会因为所有流都有了初始值,在订阅时就会立即发出包含地址信息的对象,支持仅用地址保存用户。
问题3:基于f1$、f2$、f3$的输出对mainObj对象进行操作
这个需求其实可以和问题1、2结合,核心是以mainObj为基础,累积更新三个流的输出值。除了在combineLatest里直接合并,还可以用scan操作符更优雅地处理状态累积:
import { merge } from 'rxjs'; import { scan, map } from 'rxjs/operators'; // 先把每个流映射成对应的更新函数 const firstNameUpdates$ = f1$.pipe(map(f => (obj: any) => ({...obj, firstName: f}))); const lastNameUpdates$ = f2$.pipe(map(l => (obj: any) => ({...obj, lastName: l}))); const middleNameUpdates$ = f3$.pipe(map(m => (obj: any) => ({...obj, middleName: m}))); // 合并所有更新流,用scan累积状态 merge(firstNameUpdates$, lastNameUpdates$, middleNameUpdates$) .pipe( // 初始状态用mainObj的值 scan((currentObj, updateFn) => updateFn(currentObj), (await mainObj$.toPromise())), flatMap(data => this.someHttpService.save(data)), pluck('id') ) .subscribe(id => console.log(`user ${id} saved successfully`));
这种方式更符合响应式的状态管理思路,每次只有变化的属性会被更新。
问题4:combineLatest时识别触发的Observable,仅执行对应操作
combineLatest本身会在任何流触发时传递所有流的最新值,但如果只想执行触发流对应的操作,我们可以给每个流打上标识,再用merge+scan来处理:
完整代码示例
import { merge, of } from 'rxjs'; import { scan, map, startWith, switchMap } from 'rxjs/operators'; // 给每个流添加类型标识 const f1$ = Observable.create(ob => this.eventService.on('firstName', data => ob.next({type: 'firstName', value: data}))) .pipe(startWith({type: 'firstName', value: null})); const f2$ = Observable.create(ob => this.eventService.on('lastName', data => ob.next({type: 'lastName', value: data}))) .pipe(startWith({type: 'lastName', value: null})); const f3$ = Observable.create(ob => this.eventService.on('middleName', data => ob.next({type: 'middleName', value: data}))) .pipe(startWith({type: 'middleName', value: null})); const mainObj$ = of({ address: 'abc street' }); // 先获取初始的mainObj,再处理后续更新 mainObj$.pipe( switchMap(initialObj => merge(f1$, f2$, f3$).pipe( scan((currentObj, update) => { // 根据触发的流类型,只执行对应的更新 switch(update.type) { case 'firstName': return addFirstName(update.value)(currentObj); case 'lastName': return addLastName(update.value)(currentObj); case 'middleName': return addMiddleName(update.value)(currentObj); default: return currentObj; } }, initialObj) ) ), flatMap(data => this.someHttpService.save(data)), pluck('id') ).subscribe(id => console.log(`user ${id} saved successfully`));
更简洁的写法(绑定更新函数)
如果觉得switch-case麻烦,可以直接把每个流和对应的更新函数绑定:
const updateStreams = [ f1$.pipe(map(val => (obj: any) => addFirstName(val)(obj))), f2$.pipe(map(val => (obj: any) => addLastName(val)(obj))), f3$.pipe(map(val => (obj: any) => addMiddleName(val)(obj))) ]; mainObj$.pipe( switchMap(initialObj => merge(...updateStreams).pipe( scan((acc, updateFn) => updateFn(acc), initialObj) ) ), flatMap(data => this.someHttpService.save(data)), pluck('id') ).subscribe(...);
这样每次只有触发的流对应的更新函数会被执行,避免了不必要的重复操作。
内容的提问来源于stack exchange,提问作者vito
相关产品推荐
相关产品推荐

