使用cyclops-react无法稳定运行并发订阅者的问题求助
解决cyclops-react并发订阅者出现重复值的问题
我来帮你梳理下这个问题哈~你遇到的重复值问题,核心原因是**FutureStream默认是冷流(Cold Stream)**——每次调用subscribe()的时候,都会重新触发整个上游流的执行(包括你写的map(v -> v-1)转换逻辑)。所以你的两个订阅者sub1和sub2各自独立触发了一次从initialStream到map的完整流程,导致每个原始值被处理了两次,自然就出现了重复的输出。
要实现多个订阅者共享同一个流的处理结果,你需要把冷流转成热流(Hot Stream),让所有订阅者接收同一批处理后的数据。在cyclops-react里,可以通过broadcast()来实现广播式的流共享,下面是修改后的代码:
ReactiveSeq<Integer> initialStream = ReactiveSeq.of(1, 2, 3, 4, 5, 6); // 先处理流,然后转成可广播的热流 ReactiveSeq<Integer> processedStream = initialStream.map(v -> v - 1); HotStream<Integer> broadcastStream = Spouts.broadcast(processedStream); ReactiveSubscriber<Integer> sub1 = Spouts.reactiveSubscriber(); ReactiveSubscriber<Integer> sub2 = Spouts.reactiveSubscriber(); // 订阅广播流 broadcastStream.subscribe(sub1); broadcastStream.subscribe(sub2); CompletableFuture future1 = CompletableFuture.runAsync(() -> sub1.reactiveStream().forEach(v -> System.out.println("1 -> " + v)) ); CompletableFuture future2 = CompletableFuture.runAsync(() -> sub2.reactiveStream().forEach(v -> System.out.println("2 -> " + v)) ); try { future1.get(); future2.get(); } catch (InterruptedException | ExecutionException e) { e.printStackTrace(); }
为什么这样修改有效?
Spouts.broadcast()会把输入的冷流包装成热流,上游的map(v -> v-1)只会执行一次,处理后的每个值会被广播给所有订阅者。- 两个订阅者现在接收的是同一批处理后的数据,不会再重复触发上游的流处理逻辑,也就不会出现重复值了。
修改后的输出大概会是这样(顺序可能因为异步执行略有不同,但每个值只会出现一次):
1 -> 0 2 -> 0 1 -> 1 2 -> 1 1 -> 2 2 -> 2 1 -> 3 2 -> 3 1 -> 4 2 -> 4 1 -> 5 2 -> 5
内容的提问来源于stack exchange,提问作者Pedro Alipio
相关产品推荐
相关产品推荐

