我有一个代码,我可以在循环中执行一段逻辑,并使用 Flux 之间有一些延迟。像这样的事情,
Flux.defer(() -> service.doSomething())
.repeatWhen(v -> Flux.interval(Duration.ofSeconds(10)))
.map(data -> mapper(data)) //map data
.takeUntil(v -> shouldContinue(v)) //checks if the loop can be terminated
.onErrorStop();
Run Code Online (Sandbox Code Playgroud)
现在,我想要增量延迟。这意味着,在前 5 分钟内,每次执行之间的延迟可能为 10 秒。然后,在接下来的 10 分钟内,延迟可以是 30 秒。并且,此后每次执行之间的延迟可以是一分钟。
我如何使用 Flux 来实现这一目标?
提前致谢。
我正在使用轮询方法定期获取数据。新数据可能随时到达。我想向我的客户公开一个反应式接口。所以,我想创建一个发布者(Flux?),它会在新数据可用时发布并通知订阅者。我怎么做?我看到的所有 Flux 示例都是针对数据已知/可用的情况。实际上,我想要类似基于队列的 Flux 之类的东西,并且我的轮询线程在找到新数据时可以继续填充队列。
我知道 aPublisher不能同时发布,但是如果我使用Flux#create(FluxSink),我可以安全地FluxSink#next同时调用吗?
换句话说,即使FluxSink#next并发调用,Spring 是否具有确保事件正确串行发布的内部魔法?
public class FluxTest {
private final Map<String, FluxSink<Item>> sinks = new ConcurrentHashMap<>();
// Store a new sink for the given ID
public void start(String id) {
Flux.create(sink -> sinks.put(id, sink));
}
// Called from different threads
public void publish(String id, Item item) {
sinks.get(id).next(item); //<----------- Is this safe??
}
}
Run Code Online (Sandbox Code Playgroud)
它的声音,我喜欢这一段的官方指南中指出,上述确实是安全的,但我不是在我的理解非常有信心。
create是一种更高级的 Flux 编程创建形式,适用于每轮多次发射,甚至来自多个线程。