相关疑难解决方法(0)

Java Flux 中的延迟增加

我有一个代码,我可以在循环中执行一段逻辑,并使用 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 来实现这一目标?

提前致谢。

java delay reactive-programming flux spring-webflux

6
推荐指数
1
解决办法
439
查看次数

如何为流数据创建 Flux/Publisher

我正在使用轮询方法定期获取数据。新数据可能随时到达。我想向我的客户公开一个反应式接口。所以,我想创建一个发布者(Flux?),它会在新数据可用时发布并通知订阅者。我怎么做?我看到的所有 Flux 示例都是针对数据已知/可用的情况。实际上,我想要类似基于队列的 Flux 之类的东西,并且我的轮询线程在找到新数据时可以继续填充队列。

java reactive-programming project-reactor

5
推荐指数
1
解决办法
2011
查看次数

从多个线程同时使用同一个 FluxSink 是否安全

我知道 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 编程创建形式,适用于每轮多次发射,甚至来自多个线程。

java concurrency spring project-reactor reactive-streams

1
推荐指数
1
解决办法
1894
查看次数