Rx.js:如何将combineLatest作为独立Observable处理WebSocket Subjects
合理实现RxJS合并双市场商品价格流的方案
先拆解你原来写法的问题:在fromArray$的map里嵌套combineLatest,会生成高阶Observable(也就是Observable<Observable<[Price, Price]>>),后续还得用mergeAll/switchAll扁平化,不仅代码冗余,逻辑也不够直观。
下面直接用combineLatest作为核心Observable来实现需求,分两种常见场景:
场景1:实时获取所有商品的双市场最新价格
只要任意商品的任意市场价格更新,就输出所有商品的双市场价格集合:
import { combineLatest } from 'rxjs'; import { map } from 'rxjs/operators'; import { webSocket } from 'rxjs/webSocket'; // 封装获取单个商品的市场价格流 const getMarketPrice$ = (marketUrl, productId) => webSocket(`${marketUrl}/prices/${productId}`); // 目标商品ID列表 const targetProducts = ['prod_001', 'prod_002', 'prod_003']; // 为每个商品生成双市场价格合并流,输出格式为 { productId, market1Price, market2Price } const productPricePairs$ = targetProducts.map(id => combineLatest([ getMarketPrice$('wss://market-a.com', id), getMarketPrice$('wss://market-b.com', id) ]).pipe( map(([priceA, priceB]) => ({ productId: id, market1Price: priceA, market2Price: priceB })) ) ); // 合并所有商品的价格流,实时输出全量最新价格 const allCombinedPrices$ = combineLatest(productPricePairs$); // 订阅使用 allCombinedPrices$.subscribe(allPrices => { console.log('全量商品双市场最新价格:', allPrices); });
场景2:单个商品价格更新时单独输出
只要某个商品的价格更新,就单独输出该商品的双市场价格:
import { merge, combineLatest } from 'rxjs'; import { map } from 'rxjs/operators'; import { webSocket } from 'rxjs/webSocket'; // 复用上面的getMarketPrice$和targetProducts const productPricePairs$ = targetProducts.map(id => combineLatest([ getMarketPrice$('wss://market-a.com', id), getMarketPrice$('wss://market-b.com', id) ]).pipe( map(([priceA, priceB]) => ({ productId: id, market1Price: priceA, market2Price: priceB })) ) ); // 合并所有商品的价格流,单个商品更新时单独输出 const individualUpdatedPrices$ = merge(...productPricePairs$); // 订阅使用 individualUpdatedPrices$.subscribe(updatedPrice => { console.log(`商品${updatedPrice.productId}价格更新:`, updatedPrice); });
核心优化点
- 把
combineLatest从map的嵌套逻辑中抽出来,直接作为每个商品的基础流生成器,逻辑更清晰 - 用
combineLatest或merge整合所有商品的流,避免高阶Observable的冗余处理 - 统一输出格式,后续业务逻辑更容易复用
内容的提问来源于stack exchange,提问作者reaver lover
相关产品推荐
相关产品推荐

