如何实现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
相关产品推荐
相关产品推荐

