- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我正在尝试弄清楚如何采用作为输入数据序列的 Flux,并行通过可能对序列重新排序的阻塞调用运行它们,然后通过第二个单线程阻塞调用。这个想法是最后的单线程调用将重新排序的并行工作输出记录到磁盘上。我正在尝试做的最终目标是并行算法是一种共识算法,它将确定数据输入的实际顺序。单线程写入是模拟按照共识算法确定的顺序处理事物。
查看this article它建议我应该将我的阻塞调用转换为在调度程序上运行的 Mono,该调度程序可以给我并行或单线程处理:
public class BlockingRemoteCall {
private final static Random r = new Random();
private final static Scheduler scheduler = Schedulers.newParallel("myWebservice", 10);
static private String blockingWebService(final String in) {
try {
// fakes blocking for up to a second
Thread.sleep((long) (1000 * r.nextFloat()));
System.out.println("webserver returned: "+in+" on "+Thread.currentThread().getName());
} catch (Exception e) {
throw new RuntimeException(e);
}
return in;
}
public static Mono<String> blockingMethodParallelThread(final String in) {
return Mono.fromSupplier(() -> blockingWebService(in))
.subscribeOn(scheduler);
}
}
public class BlockingJournal {
private final static Scheduler scheduler = Schedulers.newSingle("myJournal");
private static String blockingWrite(String in){
try {
// fakes blocking for disk write
Thread.sleep(5L);
System.out.println("journal wrote: "+in+" on "+Thread.currentThread().getName());
} catch (Exception e){
throw new RuntimeException(e);
}
return in;
}
public static Mono<String> blockingMethodSingleThread(final String in) {
return Mono.fromSupplier(() -> blockingWrite(in))
.subscribeOn(scheduler);
}
}
我一直在尝试获取整数 Flux 并通过这些方法以某种方式映射或平面映射,但我无法记录任何内容。这是我最近的尝试:
final Scheduler parallelScheduler = Schedulers.newParallel("p");
final Scheduler singleScheduler = Schedulers.single();
Flux<String> flux = Flux.range(1, 10).map(i -> i.toString()).publishOn(parallelScheduler);
Flux<String> pipeline = flux.map(s->{
Mono<String> async = BlockingRemoteCall.blockingMethodParallelThread(s);
String r1 = async.block();
Mono<String> r2 = BlockingJournal.blockingMethodSingleThread(r1);
return r2.block();
});
pipeline.subscribeOn(singleScheduler).doOnNext(System.out::println).blockLast();
这实际上并没有输出任何东西,但每当我能够生成任何输出时,我只看到 println
语句显示数据流在一个线程上按顺序处理。我希望看到的是调用 blockingMethodParallelThread(s)
时的任意延迟导致输入序列被乱序记录。
我如何设置才能并行处理输入的 Flux(最终从 reactor-netty 输入冒泡),重新排序,然后最终按顺序处理,保留重新排序?重新排序是由于并行进行阻塞调用造成的?谢谢!
最佳答案
这里有几点:
publishOn(parallelScheduler)
不会使您的 Flux
并行执行,它只意味着您的顺序 Flux
现在将发布在一个并行调度器。相反,您可能希望调用 parallel()
使其并行,然后使用您选择的调度程序指定 runOn()
。 (类似地,在并行 Flux
上调用 sequential()
将使其再次顺序。)map()
调用然后在其中阻塞没有多大意义 - 您也可以使用 flatMap()
并只返回生成的发布者。<因此,考虑到这些要点,您的代码将变为:
ParallelFlux<String> flux = Flux.range(1, 10).map(i -> i.toString()).parallel().runOn(Schedulers.elastic());
ParallelFlux<String> pipeline = flux.flatMap(s -> {
Mono<String> async = BlockingRemoteCall.blockingMethodParallelThread(s);
String r1 = async.block();
return BlockingJournal.blockingMethodSingleThread(r1);
});
pipeline.sequential().doOnNext(System.out::println).blockLast();
...如您所料,这将乱序输出结果。
关于java - 使用 Reactor 乱序处理输入通量,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/58708240/
我正在创建一个“杀死”人的命令。我希望机器人返回消息“哈!你以为!@Author 死了!”如果他们 ping 机器人。 (我如何让机器人查看它是否被 ping 过?)答案已更新并且现在可以正常工作。
我有一个在heroku 上运行的应用程序,例如my-app.herokuapp.com。但是,如果我输入 ping -c 10 my-app.herokuapp.com 在Mac终端中,它显示请求超时
我在 minikube 集群中有一个 k8s 服务/部署(default 命名空间中的名称 amq: D20181472:argo-k8s gms$ kubectl get svc --all-nam
我有 2 个 EC2 Ubuntu 实例。它们共享相同的 VPC、子网和安全组。实例的防火墙已关闭。但是私网IP还是无法互相ping通。如何让这些实例互相 ping 通? 最佳答案 在安全组中,为“回
我可以连接到我的 wifi(另一台笔记本电脑在此网络上正常),但是浏览器不会加载网页,并且我无法 ping 通 google.com 我注意到的一件奇怪的事情是,如果我查看/etc/resolv.co
我在 Azure 上使用 PUBSUB 时遇到问题。 Azure 防火墙将关闭闲置任意时间的连接。对于时间长度存在很多争议,但人们认为大约是 5 - 15 分钟。 我使用 Redis 作为消息队列。为
我很无聊,因为我的开发服务器已关闭,我正在运行命令提示符以无限期地 ping 服务器,以便我看到它们何时停止超时并知道我可以再次工作。与此同时,我想制作一个 Air 应用程序来为我做这件事,所以当它开
是否可以向 nat 后面的主机发送回显请求 后。所有的 echo-request 都不包含目标主机的端口,因此如果有多个主机使用相同的外部 ip 地址,nat 将如何将 echo-reques
我按照以下链接创建了 azure 实例 http://michaelwasham.com/2013/09/03/connecting-clouds-site-to-site-aws-azure/ 我可
friend 们,我认为这是一件奇怪的事情(至少对我来说)。因为我了解到互联网上的每个域名都有一个对应的IP地址。它存储在 DNS 上的某个位置。 现在,这就是我从命令行 ping google.co
我正在尝试使用分配给 kube-dns 服务的集群 IP 从 dnstools pod ping kube-dns 服务。 ping 请求超时。在同一个 dnstools pod 中,我尝试使用暴露的
我按照以下链接创建了 azure 实例 http://michaelwasham.com/2013/09/03/connecting-clouds-site-to-site-aws-azure/ 我可
我有一个虚拟网络 vmnet2,使用 10.0.2.0/24 网络,我希望我的 Linux 服务器能够 ping 默认网关。 我已将 Linux eth1 值设置为 IPADDR="10.0.2.50
我想将我的本地 mysql 数据库迁移到 Amazon RDS。但首先我想测试它是否正在接收通信。所以我尝试ping它。但是尝试超时。 ping -c 5 myfishdb.blackOut.us-w
我对 AWS 很陌生,已经测试过启动一个实例,如下所示: 下面是安全组,附加了inbound规则 我的问题是我无法 ping 通这台服务器。我可以知道我是否理解错了什么吗? 最佳答案 您需要为其创建新
我对 AWS 很陌生,已经测试过启动一个实例,如下所示: 下面是安全组,附加了inbound规则 我的问题是我无法 ping 通这台服务器。我可以知道我是否理解错了什么吗? 最佳答案 您需要为其创建新
如何确定 IP 地址是否可 ping 通?另外,如何使用 perl 脚本找到可 ping 的 IP 是静态的还是动态的? 最佳答案 看看 Net::Ping模块; #!/usr/bin/env per
我已经研究这个有一段时间了。对于网站 static.etreeblog.com,如果网站离线,我想更改 duv 的类。 我研究过的方法: - 使用带有图像的 onerror 标签来运行函数。-问题:我
我正在使用 OpenvSwitch-2.5.2 在两个虚拟机上设置第 2 层网络,如上图所示。 在阅读了 ovs 官方教程和其他一些文章后,我在每个虚拟机上尝试了以下命令: # on vm1 ip l
我有一个名为 backend 的 Docker 容器,它公开了一个端口 8200,并在其中的 gunicorn 后面运行了一个 django 服务器。这是我的 Dockerfile: FROM deb
我是一名优秀的程序员,十分优秀!