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

Flutter Riverpod同屏多Consumer监听不同流超时问题排查

问题

我有一个无状态组件,返回包含两个ConsumerWidget的Column。这两个组件分别监听来自两个StreamProvider的不同流,这两个StreamProvider均依赖同一个其他Provider,且返回不同类型的数据。仅显示其中一个组件时一切正常,但同时显示两个时,其中一个会在5秒后返回超时错误:Error TimeoutException after 0:00:05.000000: Future not completed(出错的组件随机,通常是第二个)。

我同时使用Flutter Blue Plus与蓝牙设备通信,以下是代码简化版本:

// responsePacket here - List<int> converted to specifically to Data1 or Data2 depending on 1st byte
responseStream = _txCharacteristic.value..map(responsePacket).where((it) => it != null).asBroadcastStream()

Stream<T> getStream<T>() => responseStream.where((element) => element is T).cast<T>();

Future<bool> _sendIt(List<int> message) {
  return _rxCharacteristic.write(requestPacket, withoutResponse: false);
}

Future<bool> requestData1() {
  return _sendIt([1]);
}

Future<bool> requestData2() {
  return _sendIt([2]);
}

final streamDataProvider1 = StreamProvider.autoDispose<Data1>((ref) async* {
  Object obj = ref.watch(objectStateNotifierProvider);

  while (obj.data != null) {
    obj.bt!.requestData1();
    Data1 value = await obj.data!.getStream<Data1>().first.timeout(const Duration(seconds: 5));

    yield value;
  }
});

final streamDataProvider2 = StreamProvider.autoDispose<Data2>((ref) async* {
  Object obj = ref.watch(objectStateNotifierProvider);

  while (obj.data != null) {
    obj.bt!.requestData2();
    Data2 value = await obj.data!.getStream<Data2>().first.timeout(const Duration(seconds: 5));

    yield value;
  }
});

class _MainScreenData extends StatelessWidget {
  @override
  Widget build(BuildContext context) {
    return Column(
      children: [
        _TestStream1(), // Comment out one of these and it works as expected
        _TestStream2(),
      ],
    );
  }
}

class _TestStream1 extends ConsumerWidget {
  @override
  Widget build(BuildContext context, WidgetRef ref) {
    final stream = ref.watch(streamDataProvider1);

    return stream.when(
      data: (value) => Text(value.toString()),
      error: (e, s) => Text('Error $e'),
      loading: () => const CircularProgressIndicator(),
    );
  }
}

class _TestStream2 extends ConsumerWidget {
  @override
  Widget build(BuildContext context, WidgetRef ref) {
    final stream = ref.watch(streamDataProvider2);

    return stream.when(
      data: (value) => Text(value.toString()),
      error: (e, s) => Text('Error $e'),
      loading: () => const CircularProgressIndicator(),
    );
  }
}

我希望能在某些场景下单独监听这些流,若在一个组件中使用两个流,我了解可以嵌套StreamBuilder并使用单独快照(尚未尝试)。请问导致其中一个Widget的Future无法完成的原因是什么?

原因分析

问题核心在于并发蓝牙请求的竞争与响应匹配混乱:

  • BLE通信的半双工特性:蓝牙BLE设备通常无法同时处理连续的两个请求,两个Provider几乎同时调用requestData1()和requestData2(),会导致设备来不及响应其中一个请求,或者响应顺序被打乱。
  • 流订阅的匹配问题:getStream<T>().first会等待第一个匹配类型的响应,但如果两个请求的响应返回顺序和订阅顺序不匹配(比如requestData2()的响应先返回,而streamDataProvider1的订阅还在等待Data1),就会出现某个订阅长时间收不到对应类型的事件,最终触发超时。
  • 广播流的订阅时机偏差:虽然responseStream是广播流,但两个Provider的订阅启动时间几乎同步,若设备响应在其中一个订阅完成注册前到达,会导致该订阅错过事件(广播流不会丢失事件,但连续请求的响应间隔极短时仍可能出现匹配错位)。

解决方案

方案1:序列化蓝牙请求

避免同时发送请求,确保上一个请求的响应返回后再发送下一个,消除竞争:

// 合并为一个Provider,按顺序请求数据
final combinedStreamProvider = StreamProvider.autoDispose<({Data1 data1, Data2 data2})>((ref) async* {
  Object obj = ref.watch(objectStateNotifierProvider);
  while (obj.data != null) {
    // 先请求Data1并等待响应
    obj.bt!.requestData1();
    Data1 data1 = await obj.data!.getStream<Data1>().first.timeout(const Duration(seconds: 5));
    
    // 再请求Data2并等待响应
    obj.bt!.requestData2();
    Data2 data2 = await obj.data!.getStream<Data2>().first.timeout(const Duration(seconds: 5));
    
    yield (data1: data1, data2: data2);
  }
});

UI中监听合并后的流,分别展示两个数据即可。

方案2:为请求添加唯一标识

给每个请求添加唯一ID,确保响应能精准匹配到对应的请求:

// 修改请求方法,加入请求ID
Future<bool> _sendIt(List<int> message, int requestId) {
  final packet = [...message, requestId]; // 将ID加入数据包
  return _rxCharacteristic.write(packet, withoutResponse: false);
}

// 解析响应时保留请求ID
class Data1WithId {
  final Data1 data;
  final int requestId;
  Data1WithId({required this.data, required this.requestId});
}

class Data2WithId {
  final Data2 data;
  final int requestId;
  Data2WithId({required this.data, required this.requestId});
}

dynamic responsePacket(List<int> value) {
  int requestId = value.last;
  if (value[0] == 1) {
    return Data1WithId(data: Data1(...), requestId: requestId);
  } else if (value[0] == 2) {
    return Data2WithId(data: Data2(...), requestId: requestId);
  }
  return null;
}

// 修改getStream,按类型+请求ID过滤
Stream<T> getStream<T>(int requestId) => responseStream
    .where((element) => element is T && (element as dynamic).requestId == requestId)
    .cast<T>()
    .map((e) => e.data);

Provider中生成唯一ID,确保响应只被对应订阅接收:

final streamDataProvider1 = StreamProvider.autoDispose<Data1>((ref) async* {
  Object obj = ref.watch(objectStateNotifierProvider);
  int requestId = 0;
  
  while (obj.data != null) {
    requestId++;
    obj.bt!.requestData1(requestId);
    Data1 value = await obj.data!.getStream<Data1WithId>(requestId).first.timeout(const Duration(seconds: 5));
    yield value;
  }
});

方案3:使用BehaviorSubject缓存数据

如果不需要连续请求,而是获取一次数据后缓存,用BehaviorSubject保存最新数据:

// 在蓝牙服务类中添加Subject
final BehaviorSubject<Data1> _data1Subject = BehaviorSubject();
final BehaviorSubject<Data2> _data2Subject = BehaviorSubject();

// 初始化时监听响应流,分发到对应Subject
void init() {
  responseStream.listen((data) {
    if (data is Data1) {
      _data1Subject.add(data);
    } else if (data is Data2) {
      _data2Subject.add(data);
    }
  });
}

// 修改Provider监听Subject
final streamDataProvider1 = StreamProvider.autoDispose<Data1>((ref) {
  Object obj = ref.watch(objectStateNotifierProvider);
  // 首次触发请求
  obj.bt!.requestData1();
  return obj.bt!._data1Subject.stream;
});

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 23:10:43