【问题标题】:RxJS: Count changes on collection itemsRxJS:计算集合项的更改
【发布时间】:2021-04-15 07:44:37
【问题描述】:

给定一组项目,在任何项目/项目属性的特定时间间隔内跟踪/监控/计数变化的最佳方法是什么?
更具体地说,我需要能够监控变化并通知特定的动态条件。

到目前为止我的解决方案:
每次发生变化时,我都会在 ReplaySubject 中推送数据更改日志。
我订阅了这个主题并应用了一些自定义运算符,这些运算符过滤特定列并计算更改。

例如。 '如果项目 Y 的属性 X 在过去 5 秒内更改了 2 次,请通知我'

// const dataChangeLog$ = new ReplaySubject(...)

    dataChangeLog$.pipe(
      filterColumnChanges('columnName'), // filters only changes on the given column
      windowTime(5000),
      switchMap((change$) => change$.pipe(countCellChanges(2)))
    );

    const countCellChanges = (
      limit: number,
      reset = true
    ): MonoTypeOperatorFunction<DataChangedLogEntry> => {
      return (source$) =>
        defer(() => {
          const counterMap = new Map<unknown, number>();

          const getCellCounter = (dataChangeLog: DataChangedLogEntry) =>
            counterMap.get(dataChangeLog.primaryKeyValue) ?? 0;

          return source$.pipe(
            tap((dataChangeLog) => {
              let currentCounter = getCellCounter(dataChangeLog);
              counterMap.set(dataChangeLog.primaryKeyValue, ++currentCounter);
            }),
            filter((dataChangeLog) => {
              return getCellCounter(dataChangeLog) >= limit;
            }),
            // if limit is reached, reset it
            tap((dataChangeLog) => {
              if (reset) {
                counterMap.set(dataChangeLog.primaryKeyValue, 0);
              }
            })
          );
        });
    };

问题是计数器每 5 秒重置一次,无论我在此间隔内进行了多少更改。
我真正需要的是检查每一次排放是否在最后 Y 秒内至少有 X 次排放。

【问题讨论】:

  • 看看bufferTime 运算符,看看它是否能帮到你:)

标签: javascript typescript rxjs


【解决方案1】:

此操作符检查源并在 您描述的条件满足。

const xInLastY = (x, y) => (source$) => {
  const runningCount$ = source$.pipe(
    mergeMap(() => {
      const inc$ = of(1);
      const dec$ = of(-1).pipe(delay(y));
      return of(inc$, dec$).pipe(mergeAll());
    }),
    scan((acc, next) => acc + next, 0),
  );

  return runningCount$.pipe(
    map((count) => count >= x),
    distinctUntilChanged(),
    filter(x => x)
  );
};

编辑:专门解决这个位:

如果项目 Y 的属性 X 在过去 5 秒内更改了 2 次,请通知我

运算符巧妙地融入了这个用例。

const recentChangesToProperty$ = source$.pipe(
  distinctUntilKeyChanged(MY_PROPERTY),
  xInLastY(2, 5e3)
);

最后更新:

这概括了“最后 X 毫秒的条件”的情况,以满足 OP 的其他要求:

仅当更改的值是最后 x 毫秒内的最高/最低时才发出

我们不是累积一个增长/缩小的计数,而是累积一个增长/缩小的值数组。 “max”的情况在下面实现,“min”的实现很容易从这里得到。

type ArrayAction<T> = { type: 'PUSH'; value: T } | { type: 'POP' };

const arrayReducer = <T>(state: T[], action: ArrayAction<T>) => {
  switch (action.type) {
    case 'PUSH':
      return [...state, action.value];
    case 'POP':
      return state.slice(1);
    default:
      return state;
  }
};

const timeTrailingValues = (ms: number) => <T>(
  source$: Observable<T>
): Observable<T[]> => {
  const runningList$ = source$.pipe(
    mergeMap((value) => {
      const in$ = of({ type: 'PUSH', value });
      const out$ = of({ type: 'POP' }).pipe(delay(y));
      return of(in$, out$).pipe(mergeAll());
    }),
    scan(arrayReducer, [])
  );
};

const maxInLastX = (ms: number) => <T>(
  source: Observable<T>
): Observable<boolean> =>
  source$.pipe(
    timeTrailingValues(ms),
    map((xs) => {
      const [latestValue] = xs.slice(-1);
      const rest = xs.slice(0, -1);
      return rest.every((value) => value <= latestValue);
    }),
    filter(isMax => isMax)
  );

// Still want to track the running count? Combine timeTrailingValues with this
const atSize = (x: number) => <T>(
  source$: Observable<T[]>
): Observable<boolean> => {
  return source$.pipe(
    map((xs) => xs.length >= x),
    distinctUntilChanged(),
    filter((x) => x)
  );
};

【讨论】:

  • 谢谢!我现在不在工作站,所以我不能玩它,但是,如果我读错了,请纠正我,这个操作符将等待 Y 毫秒,然后如果在这段时间内至少有 X 次发射,则会发出 TRUE .如果在最后最多 Y 毫秒内至少有 X 次发射,我需要立即发射。抱歉,现在我意识到我的问题并不清楚。
  • 不,它应该按您的意愿工作。这样想 - 我们正在创建一个中间 observable delta$,它只发出 1 和 -1。对于来自源的每个发射,它立即发射 1,然后在 Y 毫秒后发射 -1。通过总结这些 1 和 -1(使用 scan),我们可以获得最近 Y 毫秒内持续更新的排放计数。
  • 我更新了结果 observable 以在满足条件时立即发出(它发出连续的真/假值)。
  • 谢谢,现在我明白了 :) 显然您的 rxjs 比我的先进得多,您对我的用例有什么提示吗?我有一个列表(不是流),其中的每个更改都被推送到 changeLog 流(ReplaySubject)中。我希望能够构建动态查询并从 changeLog 中提取信息,例如:1. 此项目/属性已更改,2. 此项目/属性已更改 X 次,3. 此项目/属性已更改 X 次最后Y 秒,4. 此属性已更改为 [在最后 Y 秒内] 的最小值/最大值等
  • 我做到了,但我没有足够的声望点 ;)
猜你喜欢
  • 2021-09-16
  • 2014-03-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-08-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多