基于RxJS实现支持合并更新的分组去重延迟输出功能
基于RxJS实现支持合并更新的分组去重延迟输出功能
看起来你已经搞定了基于ids分组的去重延迟输出,现在要扩展支持merge类型的对象合并到对应分组,我来帮你调整代码实现这个需求。
核心思路就两点:
- 让merge类型的对象能精准找到对应的目标分组
- 在分组内部累积合并所有同组对象的属性
完整修改后的代码
const { Subject, groupBy, mergeMap, scan, last, filter, timer } = require('rxjs'); const waitTime = 1000; // 维护id到分组key的映射表,让merge对象能找到对应分组 const elementToGroupKey = new Map(); const mySource = new Subject(); mySource .pipe( groupBy( (val) => { if (val.ids) { // 原逻辑:带ids的对象用排序后的ids生成唯一分组key const key = JSON.stringify(val.ids.sort()); // 把每个id和分组key绑定,供后续merge对象使用 val.ids.forEach(id => elementToGroupKey.set(id, key)); return key; } else if (val.mergeMeWith) { // 带mergeMeWith的对象:通过元素找到对应分组的key // 假设mergeMeWith的元素都属于同一个分组,取第一个元素查映射 const targetId = val.mergeMeWith[0]; return elementToGroupKey.get(targetId) || JSON.stringify(val.mergeMeWith.sort()); } // 兜底:无ids/mergeMeWith的对象归到空分组 return JSON.stringify([]); }, null, () => timer(waitTime) ), mergeMap(group$ => group$.pipe( // 用scan操作符实时累积合并组内所有对象 scan((acc, curr) => { // 初始值取组内第一个对象(带ids的原始分组对象) if (!acc) { return { ...curr }; } // 开始合并属性 const merged = { ...acc }; // 对mergeMeWith数组做去重合并 if (curr.mergeMeWith) { merged.mergeMeWith = [ ...new Set([...(merged.mergeMeWith || []), ...curr.mergeMeWith]) ]; } // 合并其他自定义属性(比如foo1、foo2),直接覆盖旧值 // 如果需要保留历史值,可以改成数组存储,根据你的业务调整 Object.entries(curr).forEach(([key, value]) => { if (key !== 'ids' && key !== 'mergeMeWith') { merged[key] = value; } }); return merged; }, null), // 只在分组等待超时后,输出最终的合并结果 last(), // 过滤掉空值,避免无有效元素的分组输出 filter(Boolean) )) ) .subscribe(data => console.log(new Date().toISOString(), "subscriber GOT:", data)); // 测试代码调整(新增merge类型测试对象) async function sendToObs(data,delay) { await new Promise(resolve => setTimeout(resolve,delay)) console.log(new Date().toISOString(), "Sending", data) mySource.next(data); } async function myMain() { await sendToObs({ ids: ['a','b'] }, 100); await sendToObs({ ids: ['c','d'] }, 100); await sendToObs({ ids: ['a','b'] }, 100); await sendToObs({ ids: ['b','a'] }, 100); await sendToObs({ ids: ['a','b'] }, 100); await sendToObs({ ids: ['e','f'] }, 100); // 新增两个merge测试对象 await sendToObs({ mergeMeWith: ['e'], foo1: '7_e' }, 100); await sendToObs({ mergeMeWith: ['f'], foo2: '42_f' }, 100); await sendToObs({ ids: ['a','b'] }, waitTime + 500); } myMain();
关键逻辑说明
分组匹配优化:
- 新增
elementToGroupKey映射表,当带ids的对象进入时,把每个id和分组key绑定,这样后续的merge对象就能通过mergeMeWith里的元素找到对应分组。 - merge类型对象通过
mergeMeWith中的元素查映射表,精准归到目标分组。
- 新增
组内属性合并:
- 用
scan操作符实时累积合并组内所有对象:保留原始的ids属性确保分组标识正确,对mergeMeWith数组做去重合并,其他自定义属性直接覆盖(如果需要保留历史值可以改成数组存储)。 - 用
last()确保只在分组等待超时后输出最终合并结果,和原逻辑的延迟输出行为保持一致。
- 用
测试输出效果
运行后你会看到订阅者输出的结果包含合并后的完整对象:
... subscriber GOT: { ids: [ 'e', 'f' ], mergeMeWith: [ 'e', 'f' ], foo1: '7_e', foo2: '42_f' } ...
这样既保留了你原来的分组去重延迟输出逻辑,又完美实现了merge对象合并到对应分组的需求~
备注:内容来源于stack exchange,提问作者cis
相关产品推荐
相关产品推荐

