【问题标题】:Webflux subscriberWebflux 订阅者
【发布时间】:2021-06-27 12:38:22
【问题描述】:

我目前面临一个关于在 switchIfEmpty 函数中保存 redis 的问题。可能与我在反应式编程方面相当新的事实有关,不幸的是我很难找到合适的例子。

但是我已经能够解决它,但我很确定有更好的方法来解决它。 这是我现在写的:

public Mono<ResponseEntity<BaseResponse<Content>>> getContent(final String contentId, final String contentType){
    return redisRepository.findByKeyAndId(REDIS_KEY_CONTENT, contentId.toString()).cast(Content.class)
               .map(contentDTO -> ResponseEntity.status(HttpStatus.OK.value())
                                                .body(new BaseResponse<>(HttpStatus.OK.value(), HttpStatus.OK.getReasonPhrase(), contentDTO)))
               //here I have to defer, otherwise It would never wait for the findByKeyAndId 
               .switchIfEmpty(Mono.defer(() -> {
                   Mono<ResponseEntity<BaseResponse<Content>>> responseMono = contentService.getContentByIdAndType(contentId, contentType);
                   
                   //so far I understood I need to consume every stream I have, in order to actually carry out the task otherwise will be there waiting for a consumer.
                   //once I get what I need from the cmsService I need to put it in cache and till now this is the only way I've been able to do it
                   responseMono.filter(response -> response.getStatusCodeValue() == HttpStatus.OK.value())
                               .flatMap(contentResponse -> redisRepository.save(REDIS_KEY_CONTENT, contentId.toString(), contentResponse.getBody().getData()))
                                        .subscribe();
                   //then I return the Object I firstly retrived thru the cmsService
                   return responseMono;
               }
    ));
}

有什么更好的方法的线索或建议吗? 提前感谢您的帮助!

【问题讨论】:

  • 如果理解正确,你想从redis中获取,如果没有找到你想从另一个服务中获取值,将它保存在缓存中,然后构建你的响应并返回它。
  • @Toerktumlare 正确!我知道代码很乱,但在某种程度上它执行了整个流程。

标签: java reactive-programming spring-webflux


【解决方案1】:

从最佳实践的角度来看,有些事情不是很好,而且这可能与您认为的不完全一样:

  • 您正在订阅自己,这通常是出现问题的明显迹象。除了特殊情况,订阅通常应该留给框架。这也是您需要 Mono.defer() 的原因 - 通常,在框架在正确的时间订阅您的发布者之前,什么都不会发生,而您自己管理该发布者的订阅生命周期。
  • 框架仍将订阅您的内部发布者,只是您返回的Mono 对其结果没有任何作用。因此,您可能会调用 contentService.getContentByIdAndType() 两次,而不仅仅是一次 - 一次是在您订阅时,一次是在框架订阅时。
  • 像这样订阅内部发布者会创建一个“即发即用”类型模型,这意味着当您的反应式方法返回时,您不知道 redis 是否真的保存了它,如果您随后访问可能会导致问题以后再依赖这个结果。
  • 与上述无关,但contentId已经是一个字符串,你不需要在上面调用toString() :-)

相反,您可以考虑在您的switchIfEmpty() 块中使用delayUntil - 如果响应代码正常,这将允许您将值保存到redis,延迟直到发生这种情况,并在发生时保留原始值完毕。代码可能看起来像这样(如果没有完整的示例,很难说这是否完全正确,但它应该会给你一个想法):

return redisRepository
        .findByKeyAndId(REDIS_KEY_CONTENT, contentId).cast(Content.class)
        .map(contentDTO -> ResponseEntity.status(HttpStatus.OK.value()).body(new BaseResponse<>(HttpStatus.OK.value(), HttpStatus.OK.getReasonPhrase(), contentDTO)))
        .switchIfEmpty(
                contentService.getContentByIdAndType(contentId, contentType)
                        .delayUntil(response -> response.getStatusCodeValue() == HttpStatus.OK.value() ?
                                redisRepository.save(REDIS_KEY_CONTENT, contentId, contentResponse.getBody().getData()) :
                                Mono.empty())
        );

【讨论】:

  • 非常感谢!这不仅仅是简单的答案,而是关于一些反应式编程以及不应该做什么的正确课程。我刚刚尝试过,它可以工作:) 当然我看起来好多了,而且它绝对是可读的。感恩米勒!
  • @Maxuel 这就是计划!很高兴它成功了。不客气。
  • 关于你写的 4 点:1)“订阅通常应该留给框架。”您的意思是例如由其他人(例如邮递员)调用的 RestController 2)“调用两次”,是因为我放置了 Mono 然后返回它而实际上没有做任何事情吗?如果是这样,在返回之前应用于 Mono 的过滤器和平面图怎么样? 3)“你不知道redis是否真的保存了它”,我不明白。我可以告诉你的是,redis 总是保存它(我通过监视器看到它,并且发现能够检索它)4)一个凌乱的复制粘贴 xD
  • @Maxuel 这个例子中的框架是 Webflux,只要你通过 Postman 或类似的方式向它发出请求,它就会订阅(启动反应链)。至于你的第二点 - 请记住 Mono 是不可变的,因此我们使用反应链而不是设置器来改变其状态。因此,过滤器和平面图仅适用于您在示例中明确订阅的发布者 - 它们不会影响您返回的发布者。
  • 水晶般清澈?再次感谢您的时间和耐心??
猜你喜欢
  • 2018-09-21
  • 2018-06-11
  • 1970-01-01
  • 2020-05-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多