【问题标题】:How to create a Mono triggered by completion of a flux如何创建由完成通量触发的 Mono
【发布时间】:2020-02-12 17:33:36
【问题描述】:

我有一个从 Flux 源填充 Map 的类。我希望能够即时 updateTheData(),但我希望 getTheData() 能够返回解析为当前或待处理数据的 Mono。

基本上如果 theData != null 并且没有正在进行的通量,则返回 Mono.just(theData) 否则返回最终会发出数据的单声道。

编辑:这是迄今为止我能做的最好的事情

    class SomeClass {
        private Mono<Map<String, Data>> theData;
        private final SomeFluxService someFluxService = new SomeFluxService();

        public SomeClass() {
            updateTheData();
        }

        public void updateTheData() {
            someFluxService.get()
                .collectMap(Data::getId, Function.identity())
                .subscribe(d -> this.theData = Mono.just(d));
        }

        public Mono<Map<String, Data>> getTheData() {
            return this.theData;
        }
    }

但是还是存在updateTheData()第一次完成之前getTheData()会返回null的问题

帮助表示赞赏!

【问题讨论】:

    标签: java spring-webflux project-reactor reactor


    【解决方案1】:

    但是还是存在updateTheData()第一次完成之前getTheData()会返回null的问题

    这是因为您在订阅方法中使用消费者来仅在发出数据时更新Mono。这有点奇怪,最好简单地 cache()Mono 然后立即将其分配给您的字段:

    public void updateTheData() {
        theData = someFluxService.get()
                .collectMap(Data::getId, Function.identity())
                .cache();
    }
    

    进一步考虑,如果您的要求是每x 秒更新一次数据,而不是固有地按需更新,您可以将其传递到缓存函数并完全取消单独的方法:

    public NewClass() {
        theData = someFluxService.get()
                .collectMap(Data::getId, Function.identity())
                .cache(Duration.ofMinutes(5));
    }
    
    public Mono<Map<String, Data>> getTheData() {
        return theData;
    }
    

    【讨论】:

    • 这看起来好多了,谢谢!但是有一个细微的区别,原始数据将继续返回当前数据,直到成功获得新数据为止。但我认为我可以通过保留 2 个单声道并交换来将其应用于您的解决方案......
    • @Brian 不同,但您的简报是仅当有可用的缓存数据并且通量不在进行中时才返回缓存数据,这就是这样做的:- )
    猜你喜欢
    • 2020-03-29
    • 2019-11-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-10-03
    • 1970-01-01
    • 2012-11-15
    • 1970-01-01
    相关资源
    最近更新 更多