You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 09:55:33