RxJS实现按唯一键重放各类型事件最新消息的方法
实现新订阅时重放各类型事件的最新值
需求说明
我有一个包含多种类型事件的流,希望新订阅者订阅时,先收到每种类型(按唯一标识分组)的最新事件,之后再接收后续产生的新事件。
- 事件类型示例:
DoorStateChange(door_id, new_state)、LockStateChange(lock_id, new_state)、Booted(boot_time) - 分组规则:通过生成去重字符串区分不同分组,比如
DoorStateChange/${door_id}、Booted(这类全局事件唯一标识就是自身类型)
示例场景
已收到的事件序列:
Booted(129048124093) DoorStateChange(door1, Opened) DoorStateChange(door2, Closed) DoorStateChange(door1, Closed) LockStateChange(door1, Locked) LockStateChange(door2, Locked) LockStateChange(door1, Opened)
新订阅者应立即收到以下内容,之后再接收后续新事件:
Booted(129048124093) DoorStateChange(door2, Closed) DoorStateChange(door1, Closed) LockStateChange(door2, Locked) LockStateChange(door1, Opened)
尝试的代码及问题
我尝试了以下RxJS代码,但无法实现每个分组的最新消息重放:
const { interval, of, timeout, ReplaySubject } = Rx; const { groupBy, mergeMap, share, switchMap, delay, shareReplay } = RxOperators; const int = interval(300).pipe(share()); const obs = int.pipe( share(), groupBy(x => x % 3, undefined, undefined, () => new ReplaySubject(1)), mergeMap(g => g.pipe(shareReplay(1))), ); const allEvents = obs; const subscribedLater = of(1).pipe(delay(2350),switchMap(() => obs)); of(allEvents, subscribedLater)
问题:使用ReplaySubject或shareReplay都无法让新订阅者拿到每个分组的最后一条消息。
解决方案
核心思路是:按唯一标识分组后,每个分组单独保留最新值,再合并所有分组并共享整个流,确保新订阅时能触发所有分组的最新值重放。
正确实现代码
const { interval, of, merge } = Rx; const { groupBy, mergeMap, shareReplay, map, delay, switchMap } = RxOperators; // 模拟多类型事件流,替换成你的实际事件源 const eventStream = interval(300).pipe( map(x => { // 生成模拟事件,对应不同分组 if (x % 4 === 0) return { type: 'Booted', payload: { bootTime: Date.now() } }; if (x % 4 === 1) return { type: 'DoorStateChange', payload: { doorId: 'door1', state: x % 2 === 0 ? 'Opened' : 'Closed' } }; if (x % 4 === 2) return { type: 'DoorStateChange', payload: { doorId: 'door2', state: x % 2 === 0 ? 'Closed' : 'Opened' } }; return { type: 'LockStateChange', payload: { lockId: `door${x%2+1}`, state: x % 2 === 0 ? 'Locked' : 'Opened' } }; }), share() // 先共享原始事件流,避免重复生成 ); // 生成每个事件的唯一分组键 const getGroupKey = event => { if (event.type === 'Booted') return 'Booted'; if (event.type === 'DoorStateChange') return `DoorStateChange/${event.payload.doorId}`; if (event.type === 'LockStateChange') return `LockStateChange/${event.payload.lockId}`; return event.type; }; // 处理后的事件流:新订阅时重放所有分组的最新值 const sharedEventStream = eventStream.pipe( groupBy(event => getGroupKey(event)), mergeMap(group => group.pipe(shareReplay(1))), // 每个分组保留最新值 shareReplay({ bufferSize: Infinity, refCount: true }) // 共享整个合并后的流,保留所有分组的最新值 ); // 测试:第一个订阅者 sharedEventStream.subscribe(event => console.log('初始订阅:', JSON.stringify(event))); // 延迟订阅的新订阅者 of(null).pipe( delay(2350), switchMap(() => sharedEventStream) ).subscribe(event => console.log('延迟订阅:', JSON.stringify(event)));
关键要点
- 分组键生成:确保每个需要独立保留最新值的事件类型/实例有唯一的分组键,比如
DoorStateChange/door1和DoorStateChange/door2是两个独立分组 - 分组内重放:每个分组用
shareReplay(1),确保该分组的新订阅能拿到最新值 - 全局流共享:合并后的流用
shareReplay({ bufferSize: Infinity, refCount: true }),这样新订阅时会一次性拿到所有现有分组的最新值,之后继续接收新事件
内容的提问来源于stack exchange,提问作者Ákos Vandra-Meyer
相关产品推荐
相关产品推荐

