如何让RxJS v7的所有流默认变为热流(无需每次调用share)
问题描述
我正在使用RxJS v7,发现当从零创建新的Observable,或使用map、scan、merge等操作符生成Observable时,若多次引用同一个Observable变量,它会被多次执行,除非使用share操作符。示例代码及运行结果如下:
示例代码
let share_a = false let share_b = false if (process.argv[2] === 'share_a') { share_a = true } if (process.argv[3] === 'share_b') { share_b = true } console.log({ share_a, share_b }) const { map, merge, tap, share, Observable } = require('rxjs') const sleep = ms => new Promise(resolve => setTimeout(resolve, ms)) let a = new Observable(async subscriber => { await sleep(100) console.log('a called') subscriber.next('a') }) if (share_a) { a = a.pipe(share()) } let b = a.pipe(map(() => 'b'), tap(x => console.log(`${x} called`))) if (share_b) { b = b.pipe(share()) } const c = merge(a, b).pipe(map(() => 'c'), tap(x => console.log(`${x} called`))) merge(a,b,c).subscribe(console.log)
运行结果
% node src/rxjsTest.js { share_a: false, share_b: false } a called a a called b called b a called c called c a called b called c called c % node src/rxjsTest.js share_a { share_a: true, share_b: false } a called a b called b c called c b called c called c % node src/rxjsTest.js share_a share_b { share_a: true, share_b: true } a called a b called b c called c c called c %
如何无需每次调用share操作符,就能让所有Observable实现共享?
解决方案
RxJS 里默认创建的都是冷Observable,这类Observable的核心特点就是:每一次订阅都会从头触发一遍数据源的执行逻辑——这就是你看到多次引用同一Observable时,它反复执行的根本原因。share() 本质是通过 publish()(把冷Observable转成可连接的热Observable)加 refCount()(自动管理订阅计数)实现多订阅者共享同一数据源的转换。
如果不想每次手动调用share(),可以用以下几种实用方式:
1. 封装自定义创建函数,默认生成共享Observable
自己写一个工具函数替代原生的new Observable(),内部自动帮你加上share():
const { Observable, share } = require('rxjs'); // 自定义创建共享Observable的函数 function createShared(subscribeFn) { return new Observable(subscribeFn).pipe(share()); } // 用这个函数替代new Observable() let a = createShared(async subscriber => { await sleep(100); console.log('a called'); subscriber.next('a'); });
2. 封装常用操作符,自动附加共享逻辑
如果希望map、merge这类操作符生成的Observable也默认共享,可以封装这些操作符,让它们返回的Observable自动应用share():
const { map: origMap, merge: origMerge, tap, share } = require('rxjs'); // 封装map操作符 function map(project) { return source => origMap(project)(source).pipe(share()); } // 封装merge操作符 function merge(...obs) { return origMerge(...obs).pipe(share()); } // 使用封装后的操作符生成的Observable默认就是共享的 let b = a.pipe(map(() => 'b'), tap(x => console.log(`${x} called`))); const c = merge(a, b).pipe(map(() => 'c'), tap(x => console.log(`${x} called`)));
重要提醒:别全局修改RxJS默认行为
虽然理论上可以通过修改Observable原型、全局替换原生操作符来实现“全局默认共享”,但这种做法会彻底破坏RxJS的原生语义——RxJS设计冷Observable为默认是有合理性的,全局修改后,团队协作或引入第三方依赖时很容易出现莫名其妙的bug,绝对不推荐这么做。
内容的提问来源于stack exchange,提问作者pandora2000
相关产品推荐
相关产品推荐

