【问题标题】:Spring-WebFlux Flux fails with ContextSpring-WebFlux Flux 因上下文而失败
【发布时间】:2019-08-29 01:20:54
【问题描述】:

我想在我的 Flux 管道中使用 Context 来绕过过滤。

这是我所拥有的:

public Flux<Bar> realtime(Flux<OHLCIntf> ohlcIntfFlux) {
        return Flux.zip(
                ohlcIntfFlux,
                ohlcIntfFlux.skip(1),
                Mono.subscriberContext().map(c -> c.getOrDefault("isRealtime", false))
        )
                .filter(l ->
                        l.getT3() ||
                        (!l.getT2().getEndTimeStr().equals(l.getT1().getEndTimeStr())))
                .map(Tuple2::getT1)
                .log()
                .map(this::
}

这是对此的输入:

    public void setRealtime(Flux<Bar> input) {
        Flux.zip(input, Mono.subscriberContext())
                .doOnComplete(() -> {
...
   })
                .doOnNext(t -> {
...
    })
    .subscribe()
}

我可以告诉我... 中的代码没有失败,我什至可以访问Context 映射,但是当第一次迭代完成时,我得到:

onContextUpdate(Context1{reactor.onNextError.localStrategy=reactor.core.publisher.OnNextFailureStrategy$ResumeStrategy@35d5ac51})

订阅者断开连接。

所以我的问题是我是否正确使用它,这里有什么问题?

编辑:

当我使用它的价值时,我尝试repeat()Mono.subscriberContext()

        return Flux.zip(
                ohlcIntfFlux,
                ohlcIntfFlux.skip(1),
                Mono.subscriberContext()
                        .map(c -> c.getOrDefault("isRealtime", new AtomicBoolean())).repeat()
        )
                .filter(l ->
                        l.getT3().get() ||
                                (!l.getT2().getEndTime().isEqual(l.getT1().getEndTime())))
                .map(Tuple2::getT1)

并将AtomicBoolean 设置为订阅者端的上下文,并在我需要上游信号时更改此变量引用中的值,但它根本没有改变:

        input
                .onErrorContinue((throwable, o) -> throwable.getMessage())
                .doOnComplete(() -> {
                    System.out.println("Number of trades for the strategy: " + tradingRecord.getTradeCount());
                    // Analysis
                    System.out.println("Total profit for the strategy: " + new TotalProfitCriterion().calculate(timeSeries, tradingRecord));
                })
                .doOnNext(this::defaultRealtimeEvaluator)
                .subscriberContext(Context.of("isRealtime", isRealtimeAtomic))
                .subscribe();

至少重复Flux 不会断开连接,但我从中得到的值没有更新。我没有其他线索。

Spring-webflux:2.1.3.RELEASE

【问题讨论】:

    标签: java spring-webflux project-reactor


    【解决方案1】:

    这行得通:

            input
                    .onErrorContinue((throwable, o) -> throwable.getMessage())
                    .doOnComplete(() -> { ... }
                    .flatMap(bar -> Mono.subscriberContext()
                            .map(c -> Tuples.of(bar, c)))
                    .doOnNext(this::defaultRealtimeEvaluator)
                    .subscriberContext(Context.of("isRealtime", new AtomicBoolean()))
                    .subscribe();
    

    所以重点是在我的情况下将AtomicBoolean 设置为 cotnext,然后如果您想更改它的值,则将这个变量从上下文中提取出来。上游通量也一样。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-12-12
      • 2018-09-09
      • 2019-05-18
      • 2019-01-11
      • 2011-07-15
      • 2019-05-01
      • 2012-12-30
      • 2021-09-30
      相关资源
      最近更新 更多