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

如何实现RxJS中带输出大小限制的无损Buffer功能?

实现带输出大小限制的RxJS Buffer功能

需求明确:现有可观察流myInterval,需通过通知可观察流bufferBy触发缓冲区输出;每次bufferBy触发时,输出的缓冲区项数不得超过指定的maxBufferSize,且必须保证数据无损。RxJS原生buffer操作符无此类重载,需自定义实现。

实现思路

核心逻辑是维护一个内部缓存数组:

  • 当源流(myInterval)发射新值时,将值存入缓存
  • 当通知流(bufferBy)触发时,从缓存中按maxBufferSize分批取出数据并发射,直到缓存清空
  • 处理源流的错误与完成事件时,确保缓存中剩余数据全部发射,避免丢失

自定义操作符实现

import { Observable, OperatorFunction, Subscriber, Subscription } from 'rxjs';

function bufferWithMaxSize<T>(closingNotifier: Observable<any>, maxSize: number): OperatorFunction<T, T[]> {
  return (source) => {
    return new Observable((subscriber) => {
      const buffer: T[] = [];
      let sourceSub: Subscription;
      let closingSub: Subscription;

      // 处理源流的新值,存入缓存
      sourceSub = source.subscribe({
        next(value) {
          buffer.push(value);
        },
        error(err) {
          // 出错时先发射缓存剩余数据,再传递错误
          if (buffer.length > 0) {
            subscriber.next(buffer.splice(0));
          }
          subscriber.error(err);
        },
        complete() {
          // 完成时发射缓存剩余数据,再标记完成
          if (buffer.length > 0) {
            subscriber.next(buffer.splice(0));
          }
          subscriber.complete();
        }
      });

      // 处理通知流的触发事件,分批发射缓存数据
      closingSub = closingNotifier.subscribe({
        next() {
          while (buffer.length > 0) {
            const chunk = buffer.splice(0, maxSize);
            subscriber.next(chunk);
          }
        },
        error(err) {
          subscriber.error(err);
        },
        complete() {
          // 通知流完成时,发射剩余缓存数据,再完成
          if (buffer.length > 0) {
            subscriber.next(buffer.splice(0));
          }
          subscriber.complete();
        }
      });

      // 清理订阅
      return () => {
        sourceSub.unsubscribe();
        closingSub.unsubscribe();
      };
    });
  };
}

使用示例

import { interval, fromEvent } from 'rxjs';

const maxBufferSize = 5;
const myInterval = interval(1000); // 每秒发射一个递增数值
const bufferBy = fromEvent(document, 'click'); // 点击文档触发缓冲区输出

const bufferedInterval = myInterval.pipe(bufferWithMaxSize(bufferBy, maxBufferSize));

// 订阅测试
bufferedInterval.subscribe(chunk => {
  console.log(chunk);
  // 示例输出:
  // 第一次点击(假设已积累3个值):[0,1,2]
  // 第二次点击(假设又积累5个值):[3,4,5,6,7]
  // 第三次点击(假设积累2个值):[8,9]
});

关键细节说明

  • 数据无损保障:无论是源流完成、通知流完成,还是源流出错,都会先将缓存中剩余的所有数据发射出去,再处理对应事件
  • 分批逻辑:每次通知流触发时,循环从缓存中取出最多maxSize个数据,直到缓存为空,确保单次触发不会输出超过限制的项数
  • 订阅清理:在Observable销毁时,同时取消源流和通知流的订阅,避免内存泄漏

内容的提问来源于stack exchange,提问作者ket

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 01:45:36