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
相关产品推荐
相关产品推荐

