Rxjs如何将指定时长内gRPC服务的observable数据流聚合为数组
错误原因
原有尝试代码中使用的mergeMap(_ => timer(5000))会将所有上游推送的gRPC业务流数据全部替换为timer的触发值,后续scan运算符接收不到原始业务数据,因此聚合出来的数组为空。
修复方案
1. 先补全Service层流的健壮性逻辑
原有gRPC流监听缺少错误、结束回调和取消订阅时的清理逻辑,容易引发内存泄漏,修改MyTestService.ts的listCoreStream方法:
listCoreStream(): Observable<TestReply.AsObject> { return new Observable(obs => { const req = new SomeRequest(); // 从gRPC服务获取流数据 const stream = this.client.getCoreUpdates(req); stream.on('data', (message: any) => { obs.next(message.toObject() as TestReply.AsObject); }); // 监听流错误 stream.on('error', (err) => obs.error(err)); // 监听流结束 stream.on('end', () => obs.complete()); // 订阅取消时关闭gRPC流 return () => stream.cancel(); }); }
2. 按时间窗口聚合流数据
如果需要固定每X秒输出一次该时间段内收到的所有流数据数组,直接使用RxJS内置的bufferTime运算符即可,无需手动维护临时数组:
修改MyComponent.ts的ngOnInit逻辑:
import { bufferTime } from 'rxjs'; ngOnInit(): void { // 示例设置X=2000,即每2秒聚合一次该窗口内的所有流数据 this._subscription = this._MyTestService.coreData$ .pipe( bufferTime(2000) ) .subscribe((dataList: TestReply.AsObject[]) => { // 此处dataList即为2秒内收到的所有流数据数组,直接遍历处理即可 dataList.forEach(data => { if (data) { let obj = JSON.parse(data); // 原有处理逻辑:多条件判断、样式动态绑定等 } }) }); }
3. 时间+条数双条件聚合(可选)
如果需要同时满足「最多收满Y条就输出,不用等时间窗口结束」的规则,可组合运算符实现双触发条件:
import { buffer, race, timer, skip, take } from 'rxjs'; ngOnInit(): void { // 示例设置:2秒/满5条 两个条件哪个先触发就输出当前聚合的数组 this._subscription = this._MyTestService.coreData$ .pipe( buffer( race( timer(2000), // 2秒时间到触发 this._MyTestService.coreData$.pipe(skip(4), take(1)) // 收满5条触发(索引从0计数) ) ) ) .subscribe((dataList: TestReply.AsObject[]) => { // 处理聚合后的数组逻辑 }); }
内容的提问来源于stack exchange,提问作者mv_05
相关产品推荐
相关产品推荐

