【问题标题】:Usage of elapsed() function work on Mono?在 Mono 上使用 elapsed() 函数?
【发布时间】:2019-11-19 23:57:48
【问题描述】:

我正在尝试在反应式编程中获取从 redis 读取的执行时间,在查找文档时我可以看到 elapsed() 方法将执行相同的操作并实现如下代码。

Flux.fromIterable(getActions(httpHeaders))
                .parallel()
                .runOn(Schedulers.parallel())
                .flatMap(actionFact -> methodToReadFromCache(actionFact))
                .sequential();

public Mono<ActionFact> methodToReadFromCache(actionFact) {
    return Mono.fromCallable(() -> getKey(actionFact))
                .flatMap(cacheKey ->
                  redisOperations.hasKey(key)
                                .flatMap(aBoolean -> {
                                    if (aBoolean) {
                                        return redisOperations.opsForValue().get(cacheKey);
                                    }
                                    return authzService.getRolePermissions(actionFact)
                                            .flatMap(policySetResponse ->
                                                    //save in cache
                                            );
                                })
                                .elapsed()
                                .flatMap(lambda -> {
                                    LOG.info("cache/service processing key:{}, time:{}", key, lambda.getT1());
                                    return Mono.just(lambda.getT2());
                                });

输出:

cache/service processing key:KEY1, time:3 
cache/service processing key:KEY2, time:4 
cache/service processing key:KEY3, time:18 
cache/service processing key:KEY4, time:34 
cache/service processing key:KEY5, time:46 
cache/service processing key:KEY6, time:57 
cache/service processing key:KEY7, time:70 
cache/service processing key:KEY8, time:81 
cache/service processing key:KEY9, time:91 
cache/service processing key:KEY10, time:103
cache/service processing key:KEY11, time:112
cache/service processing key:KEY12, time:121
cache/service processing key:KEY13, time:134
cache/service processing key:KEY14, time:146
cache/service processing key:KEY15, time:159

我预计每个缓存请求所花费的时间将像第一个和第二个请求一样小于 5 毫秒,但情况并非如此。 elapsed() 是否将当前获取时间添加到累积中?根据我的理解,从通量发出的每个项目都是独立的?

【问题讨论】:

    标签: java spring-webflux project-reactor spring-data-redis-reactive


    【解决方案1】:

    Mono#elapsed 测量从订阅 MonoMono 发出项目 (onNext) 之间的时间。

    在您的情况下,导致订阅和计时器启动的原因是调用 methodToReadFromCache 的外部并行化 flatMap

    导致 onNext 并因此计时的是 hasKey 和 if/else 部分的组合(redisOperations.opsForValue().get(cacheKey)authzService)。

    由于我们处于并行模式,因此外部 flatMap 的计时器数量应至少与 CPU 数量一样多。

    但时间偏差的事实可能暗示了某些事情正在阻塞或容量有限的事实。比如redisTemplate一次只能处理几个key?

    【讨论】:

    • 我不确定他的代码示例,但请注意他的 methodToReadFromCache 方法返回 void,但他有一个 return 语句。代码示例有点偏离。
    • 谢谢@Simon Basie,您能否解释一下“由于我们处于并行模式,外部 flatMap 至少应该有与 CPU 一样多的计时器”?顺便说一句,我正在使用 reactiveRedisTemplate 所以它不应该阻塞
    • @ThomasAndolf 更新了返回类型,我是从记事本复制的,所以错过了。
    • @Simon,我观察到了另一种行为,例如即使并行运行,同时从缓存中读取仅使用一个线程,但是当我关闭缓存时,我可以看到多个线程正在用于执行服务来电。由于 Redis 是单线程服务器,webflux 不能使用多线程?
    • @akreddy.21 我的意思是你使用parallel().runOn(Schedulers.parallel()) 所以它应该在所有CPU之间分配工作。假设你有 8 个 CPU,这应该会导致 8 个定时 Monos 并行,所以 8 个elapsed 持续时间应该相等或接近。如果将调度程序替换为Schedulers.newParallel(),它的行为如何?如果您在最后一个 flatMap(记录日志的那个)之前执行 publishOn(Schedulers.newParallel()) 怎么办?
    【解决方案2】:

    根据文档

    我想将排放与测量的时间(Tuple2&lt;Long, T&gt;) 联系起来……​

    • 订阅后:已过期

    • 从时间的黎明开始(嗯,计算机时间):时间戳

    elapsed 是自订阅以来的测量时间。因此,您订阅并开始发送,自您订阅服务以来的时间会增加。

    official docs

    【讨论】:

    • 请注意,对于Flux,参考指南的这一部分有点误导。只有第一个onNext 是根据订阅时间来衡量的。以下onNext 时间在 N 和 N-1 之间。但在这里我们处理的是Mono,所以这是正确的。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-07-21
    • 1970-01-01
    • 2020-12-28
    • 1970-01-01
    相关资源
    最近更新 更多