如何解决RxDart shareReplay不是“热门”吗?
我在应用程序中无法成功使用shareReplay
(0.24.1)中的RxDart
。我有一个Stream
,我想“缓存”最新发出的值,以便任何“迟到”的侦听器都可以立即获取它(流可以闲置更长的时间)。我无法执行此操作,因此我开始在应用程序外部进行实验,必须承认我不完全了解正在发生的事情。
这是一些代码:
import 'package:rxdart/rxdart.dart';
void main() async {
final s = Stream.fromIterable([1,2,3]).shareReplay(maxSize: 10);
print('First:');
await for (final i in s) {
print(i);
}
print('Second:');
await for (final i in s) {
print(i);
}
}
我希望代码能够打印
First
1
2
3
Second
1
2
3
,但是它永远不会到达第二个await for
,甚至不输出字符串Second:
;但是,程序成功完成,因此第一个await for
确实完成了(因为原始Stream
是有限的)。发生了什么事?
实际上,我拥有的Stream
可能在事件之间很长一段时间内是无尽的),并且我只对最后一个事件感兴趣,因此这里有一些模拟此的代码:
import 'package:rxdart/rxdart.dart';
void main() async {
final s = Stream.periodic(Duration(seconds: 1),(e) => e).shareReplay(maxSize: 1);
// print(s.isBroadcast);
await Future.delayed(Duration(seconds: 3));
print('First');
await for (final i in s.take(1)) {
print(i);
}
print('Second');
await for (final i in s.take(1)) {
print(i);
}
}
为了确保Stream
可以完成(因此await for
可以完成),我使用take(1)
。我希望输出为:
First
2
Second
2
(也许二三分?) 我期望发生的事情是这样:
- 定期流开始每秒发射一次(广播/热点)。
- 由于
Future.delayed
,应用程序等待3秒钟。 - 第一个迟到的侦听器出现并获得最新值2(或3),并由于
take(1)
而结束。 - 第二个延迟的侦听器出现并获得最新的值,因为在此期间流没有发出任何其他值(即使这样做,该值也将增加1,这是可以的),并且由于{ {1}}。
- 该应用程序完成。
但是,输出为:
take(1)
第一个First
0
Second
Unhandled exception:
type '_TakeStream<int>' is not a subtype of type 'Future<bool>'
#0 _StreamIterator._onData (dart:async/stream_impl.dart:1067:19)
#1 _RootZone.runUnaryGuarded (dart:async/zone.dart:1384:10)
#2 _BufferingStreamSubscription._sendData (dart:async/stream_impl.dart:357:11)
#3 _BufferingStreamSubscription._add (dart:async/stream_impl.dart:285:7)
#4 _ForwardingStreamSubscription._add (dart:async/stream_pipe.dart:127:11)
#5 _TakeStream._handleData (dart:async/stream_pipe.dart:318:12)
#6 _ForwardingStreamSubscription._handleData (dart:async/stream_pipe.dart:157:13)
#7 _RootZone.runUnaryGuarded (dart:async/zone.dart:1384:10)
#8 _BufferingStreamSubscription._sendData (dart:async/stream_impl.dart:357:11)
#9 _BufferingStreamSubscription._add (dart:async/stream_impl.dart:285:7)
#10 _SyncBroadcastStreamController._sendData (dart:async/broadcast_stream_controller.dart:385:25)
#11 _BroadcastStreamController.add (dart:async/broadcast_stream_controller.dart:250:5)
#12 _StartWithStreamSink._safeAddFirstEvent (package:rxdart/src/transformers/start_with.dart:56:12)
#13 _StartWithStreamSink.onListen (package:rxdart/src/transformers/start_with.dart:37:11)
#14 forwardStream.<anonymous closure>.<anonymous closure> (package:rxdart/src/utils/forwarding_stream.dart:31:37)
#15 forwardStream.runCatching (package:rxdart/src/utils/forwarding_stream.dart:24:12)
#16 forwardStream.<anonymous closure> (package:rxdart/src/utils/forwarding_stream.dart:31:16)
#17 _runGuarded (dart:async/stream_controller.dart:847:24)
#18 _BroadcastStreamController._subscribe (dart:async/broadcast_stream_controller.dart:213:7)
#19 _ControllerStream._createSubscription (dart:async/stream_controller.dart:860:19)
#20 _StreamImpl.listen (dart:async/stream_impl.dart:493:9)
#21 DeferStream.listen (package:rxdart/src/streams/defer.dart:37:18)
#22 StreamView.listen (dart:async/stream.dart:1871:20)
#23 new _ForwardingStreamSubscription (dart:async/stream_pipe.dart:118:10)
#24 new _StateStreamSubscription (dart:async/stream_pipe.dart:341:9)
#25 _TakeStream._createSubscription (dart:async/stream_pipe.dart:310:16)
#26 _ForwardingStream.listen (dart:async/stream_pipe.dart:83:12)
#27 _StreamIterator._initializeOrDone (dart:async/stream_impl.dart:1041:30)
#28 _StreamIterator.moveNext (dart:async/stream_impl.dart:1028:12)
#29 main (package:shopping_list/main.dart)
<asynchronous suspension>
#30 _startIsolate.<anonymous closure> (dart:isolate-patch/isolate_patch.dart:301:19)
#31 _RawReceivePortImpl._handleMessage (dart:isolate-patch/isolate_patch.dart:168:12)
和await for
实际上完成了工作,但它获得了初始值,这意味着流没有统计出任何东西,听众来之前(不是“热”)。第二个失败,但有一个例外。我想念什么?
编辑:我对take(1)
的理解可能一直是错误的,尽管我很热,因为https://pub.dev/documentation/rxdart/latest/rx/ReplaySubject-class.html说它很热。
我稍微更改了代码:
shareReplay
现在输出为:
import 'package:rxdart/rxdart.dart';
void main() async {
final s = BehaviorSubject();
s.addStream(Stream.periodic(Duration(seconds: 1),(e) => e));
// print(s.isBroadcast);
await Future.delayed(Duration(seconds: 3));
print('First');
await for (final i in s.take(1)) {
print(i);
}
print('Second');
await for (final i in s.take(1)) {
print(i);
}
}
但是程序永远不会完成...
在没有First
2
Second
2
的情况下,输出永远不会到达第二个侦听器,因此确实会产生一些效果。我不明白。
编辑2 :我想我可能知道为什么带有take(1)
的第一个示例不起作用,这是文档摘录:
它将在第一次收听时自动开始发射项目,并在没有收听者时关闭。 在我的情况下,第一个侦听器/
shareReplay
进行订阅,并且流开始运行,当循环完成时,最后一个(也是唯一的)侦听器完成,因此从await for
返回的流将关闭。不过,仍然没有解释为什么行shareReplay
无法执行。
解决方法
您的最终节目没有终止的原因是您的广播流仍然每秒都在发出事件。
没有人在听,但这是广播流,因此不在乎。它有一个计时器,每秒钟滴答作响,并且 timer 可使程序保持活动状态,即使没有人再次收听。
您也可以通过执行Timer.periodic(Duration(seconds: 1),(_) {});
使程序保持活动状态。只要计时器处于活动状态,就会发生可能。
但是,崩溃很有趣。这表明StreamIterator
类的内部逻辑存在问题。我会调查的。
不能说任何关于shareReplay
有用的东西,我对RxDart不熟悉。
我的猜测是,第一个程序遇到了问题,因为replayShared
流没有正确终止-也就是说,完成后没有发送done
事件。这可以解释该程序提早终止的原因。第一个await for
正在等待下一个事件,因此它不会结束并转到下一个循环,但是没有实时计时器或任何类型的预定事件,因此隔离器将关闭。
快速阅读BeheaviorSubject
意味着它像流控制器一样工作,因此您可能需要在addStream
完成后将其关闭。也许:
s.addStream(...).whenComplete(s.close);
版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。