使用CombineLatestStream搭配throttle处理不同频率流的预期行为
核心行为解析
首先明确两个操作符的核心逻辑:
- CombineLatestStream:始终跟踪每个输入流的最新值,只要任意一个输入流产生新值,就立即将所有流的当前最新值合并输出。
- throttleTime:在指定的时间窗口内,只保留合并流的第一个输出值,窗口内后续的所有输出都会被忽略,直到下一个窗口开启。
针对你假设的场景:
当合并流输出[A1,B1,C1]后进入节流窗口,期间流1(加速度计)产生A2,此时CombineLatestStream会生成[A2,B1,C1]但被节流忽略;如果同时流2(方向)产生B2、流3(位置)产生C2,CombineLatestStream内部会更新这两个流的最新值,此时合并后的缓存值是[A2,B2,C2],但依然不会输出。
当下一个节流窗口开启后,只要任意一个输入流产生新值(比如流1发A3),CombineLatestStream会立即合并所有流的最新值[A3,B2,C2],这个值会作为新窗口的第一个值被输出——也就是说,不会持续使用旧的B1、C1,而是会用各流的最新值。
针对你需求的优化方案
你的目标是每20ms稳定获取所有流的最新数据,当前用throttleTime的方案有缺陷:如果节流窗口内没有任何输入流触发,就不会有输出,无法保证定时获取。更合适的方案是用定时采样的方式:
方案1:用sample操作符定时采样合并流
CombineLatestStream负责跟踪所有流的最新值,然后用定时流来触发采样,确保每20ms获取一次最新合并值:
CombineLatestStream.list([ accelerometerEventStream(samplingPeriod: SensorInterval.fastestInterval), frs.RotationSensor.orientationStream, getLocationUpdates(), ]) // 每20ms采样一次合并流的最新值 .sample(Stream.periodic(const Duration(milliseconds: 20))) .listen((data) { final accEvent = data[0] as AccelerometerEvent; final frsEvent = data[1] as frs.OrientationEvent; final locEvent = data[2] as Position; // 处理数据 });
方案2:用interval结合withLatestFrom
另一种方式是定时触发,每次触发时获取各流的最新值:
Stream.periodic(const Duration(milliseconds: 20)) .withLatestFrom3( accelerometerEventStream(samplingPeriod: SensorInterval.fastestInterval), frs.RotationSensor.orientationStream, getLocationUpdates(), (_, acc, orientation, loc) => [acc, orientation, loc], ) .listen((data) { final accEvent = data[0] as AccelerometerEvent; final frsEvent = data[1] as frs.OrientationEvent; final locEvent = data[2] as Position; // 处理数据 });
这两种方案都能保证每20ms获取一次所有流的最新值,不受输入流触发频率的影响,更符合你的需求。
内容的提问来源于stack exchange,提问作者joshd
相关产品推荐
相关产品推荐

