You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于RxJS实现支持合并更新的分组去重延迟输出功能

基于RxJS实现支持合并更新的分组去重延迟输出功能

看起来你已经搞定了基于ids分组的去重延迟输出,现在要扩展支持merge类型的对象合并到对应分组,我来帮你调整代码实现这个需求。

核心思路就两点:

  1. 让merge类型的对象能精准找到对应的目标分组
  2. 在分组内部累积合并所有同组对象的属性

完整修改后的代码

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();

关键逻辑说明

  1. 分组匹配优化:

    • 新增elementToGroupKey映射表,当带ids的对象进入时,把每个id和分组key绑定,这样后续的merge对象就能通过mergeMeWith里的元素找到对应分组。
    • merge类型对象通过mergeMeWith中的元素查映射表,精准归到目标分组。
  2. 组内属性合并:

    • 用scan操作符实时累积合并组内所有对象:保留原始的ids属性确保分组标识正确,对mergeMeWith数组做去重合并,其他自定义属性直接覆盖(如果需要保留历史值可以改成数组存储)。
    • 用last()确保只在分组等待超时后输出最终合并结果,和原逻辑的延迟输出行为保持一致。

测试输出效果

运行后你会看到订阅者输出的结果包含合并后的完整对象:

...
subscriber GOT: { ids: [ 'e', 'f' ], mergeMeWith: [ 'e', 'f' ], foo1: '7_e', foo2: '42_f' }
...

这样既保留了你原来的分组去重延迟输出逻辑,又完美实现了merge对象合并到对应分组的需求~

备注:内容来源于stack exchange,提问作者cis

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 10:38:02