RxDart shareReplay不是“热门”吗?

如何解决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

(也许二三分?) 我期望发生的事情是这样:

  1. 定期流开始每秒发射一次(广播/热点)。
  2. 由于Future.delayed,应用程序等待3秒钟。
  3. 第一个迟到的侦听器出现并获得最新值2(或3),并由于take(1)而结束。
  4. 第二个延迟的侦听器出现并获得最新的值,因为在此期间流没有发出任何其他值(即使这样做,该值也将增加1,这是可以的),并且由于{ {1}}。
  5. 该应用程序完成。

但是,输出为:

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 举报,一经查实,本站将立刻删除。

相关推荐


使用本地python环境可以成功执行 import pandas as pd import matplotlib.pyplot as plt # 设置字体 plt.rcParams[&#39;font.sans-serif&#39;] = [&#39;SimHei&#39;] # 能正确显示负号 p
错误1:Request method ‘DELETE‘ not supported 错误还原:controller层有一个接口,访问该接口时报错:Request method ‘DELETE‘ not supported 错误原因:没有接收到前端传入的参数,修改为如下 参考 错误2:cannot r
错误1:启动docker镜像时报错:Error response from daemon: driver failed programming external connectivity on endpoint quirky_allen 解决方法:重启docker -&gt; systemctl r
错误1:private field ‘xxx‘ is never assigned 按Altʾnter快捷键,选择第2项 参考:https://blog.csdn.net/shi_hong_fei_hei/article/details/88814070 错误2:启动时报错,不能找到主启动类 #
报错如下,通过源不能下载,最后警告pip需升级版本 Requirement already satisfied: pip in c:\users\ychen\appdata\local\programs\python\python310\lib\site-packages (22.0.4) Coll
错误1:maven打包报错 错误还原:使用maven打包项目时报错如下 [ERROR] Failed to execute goal org.apache.maven.plugins:maven-resources-plugin:3.2.0:resources (default-resources)
错误1:服务调用时报错 服务消费者模块assess通过openFeign调用服务提供者模块hires 如下为服务提供者模块hires的控制层接口 @RestController @RequestMapping(&quot;/hires&quot;) public class FeignControl
错误1:运行项目后报如下错误 解决方案 报错2:Failed to execute goal org.apache.maven.plugins:maven-compiler-plugin:3.8.1:compile (default-compile) on project sb 解决方案:在pom.
参考 错误原因 过滤器或拦截器在生效时,redisTemplate还没有注入 解决方案:在注入容器时就生效 @Component //项目运行时就注入Spring容器 public class RedisBean { @Resource private RedisTemplate&lt;String
使用vite构建项目报错 C:\Users\ychen\work&gt;npm init @vitejs/app @vitejs/create-app is deprecated, use npm init vite instead C:\Users\ychen\AppData\Local\npm-
参考1 参考2 解决方案 # 点击安装源 协议选择 http:// 路径填写 mirrors.aliyun.com/centos/8.3.2011/BaseOS/x86_64/os URL类型 软件库URL 其他路径 # 版本 7 mirrors.aliyun.com/centos/7/os/x86
报错1 [root@slave1 data_mocker]# kafka-console-consumer.sh --bootstrap-server slave1:9092 --topic topic_db [2023-12-19 18:31:12,770] WARN [Consumer clie
错误1 # 重写数据 hive (edu)&gt; insert overwrite table dwd_trade_cart_add_inc &gt; select data.id, &gt; data.user_id, &gt; data.course_id, &gt; date_format(
错误1 hive (edu)&gt; insert into huanhuan values(1,&#39;haoge&#39;); Query ID = root_20240110071417_fe1517ad-3607-41f4-bdcf-d00b98ac443e Total jobs = 1
报错1:执行到如下就不执行了,没有显示Successfully registered new MBean. [root@slave1 bin]# /usr/local/software/flume-1.9.0/bin/flume-ng agent -n a1 -c /usr/local/softwa
虚拟及没有启动任何服务器查看jps会显示jps,如果没有显示任何东西 [root@slave2 ~]# jps 9647 Jps 解决方案 # 进入/tmp查看 [root@slave1 dfs]# cd /tmp [root@slave1 tmp]# ll 总用量 48 drwxr-xr-x. 2
报错1 hive&gt; show databases; OK Failed with exception java.io.IOException:java.lang.RuntimeException: Error in configuring object Time taken: 0.474 se
报错1 [root@localhost ~]# vim -bash: vim: 未找到命令 安装vim yum -y install vim* # 查看是否安装成功 [root@hadoop01 hadoop]# rpm -qa |grep vim vim-X11-7.4.629-8.el7_9.x
修改hadoop配置 vi /usr/local/software/hadoop-2.9.2/etc/hadoop/yarn-site.xml # 添加如下 &lt;configuration&gt; &lt;property&gt; &lt;name&gt;yarn.nodemanager.res