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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 17:16:01