problem: CompoundReader makes unnecessary requests
This commit is contained in:
@@ -23,7 +23,8 @@ import reactor.core.publisher.Mono
|
|||||||
import java.time.Duration
|
import java.time.Duration
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Composition of multiple readers. Reader returns first value returned by any of the source readers.
|
* Composition of multiple readers.
|
||||||
|
* Reader returns first value returned by any of the source readers by checking one by one until one of them returns a non-empty result.
|
||||||
*/
|
*/
|
||||||
class CompoundReader<K, D>(
|
class CompoundReader<K, D>(
|
||||||
private vararg val readers: Reader<K, D>
|
private vararg val readers: Reader<K, D>
|
||||||
@@ -38,12 +39,13 @@ class CompoundReader<K, D>(
|
|||||||
return Mono.empty()
|
return Mono.empty()
|
||||||
}
|
}
|
||||||
return Flux.fromIterable(readers.asIterable())
|
return Flux.fromIterable(readers.asIterable())
|
||||||
.flatMap { rdr ->
|
.flatMap({ rdr ->
|
||||||
rdr.read(key)
|
rdr.read(key)
|
||||||
.timeout(Defaults.timeoutInternal, Mono.empty())
|
.timeout(Defaults.timeoutInternal, Mono.empty())
|
||||||
.doOnError { t -> log.warn("Failed to read from $rdr", t) }
|
.doOnError { t -> log.warn("Failed to read from $rdr", t) }
|
||||||
.onErrorResume { Mono.empty() }
|
.onErrorResume { Mono.empty() }
|
||||||
}.next()
|
}, 1)
|
||||||
|
.next()
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
@@ -72,21 +72,20 @@ class CompoundReaderSpec extends Specification {
|
|||||||
.verify(Duration.ofSeconds(1))
|
.verify(Duration.ofSeconds(1))
|
||||||
}
|
}
|
||||||
|
|
||||||
def "Return second"() {
|
def "Doesn't call others after getting first"() {
|
||||||
setup:
|
setup:
|
||||||
def reader = new CompoundReader<String, String>(reader3, reader2)
|
def call2 = false
|
||||||
when:
|
def reader2 = new Reader<String, String>() {
|
||||||
def act = reader.read("test")
|
@Override
|
||||||
then:
|
Mono<String> read(String key) {
|
||||||
StepVerifier.create(act)
|
call2 = true
|
||||||
.expectNext("test-2")
|
return Mono.just("test-2").delaySubscription(Duration.ofMillis(200))
|
||||||
.expectComplete()
|
}
|
||||||
.verify(Duration.ofSeconds(1))
|
}
|
||||||
}
|
def reader = new CompoundReader<String, String>(
|
||||||
|
reader1,
|
||||||
def "Return third"() {
|
reader2
|
||||||
setup:
|
)
|
||||||
def reader = new CompoundReader<String, String>(reader3, reader2, reader1)
|
|
||||||
when:
|
when:
|
||||||
def act = reader.read("test")
|
def act = reader.read("test")
|
||||||
then:
|
then:
|
||||||
@@ -94,16 +93,29 @@ class CompoundReaderSpec extends Specification {
|
|||||||
.expectNext("test-1")
|
.expectNext("test-1")
|
||||||
.expectComplete()
|
.expectComplete()
|
||||||
.verify(Duration.ofSeconds(1))
|
.verify(Duration.ofSeconds(1))
|
||||||
|
!call2
|
||||||
}
|
}
|
||||||
|
|
||||||
def "Ignore empty"() {
|
def "Return first even if it's slow"() {
|
||||||
setup:
|
setup:
|
||||||
def reader = new CompoundReader<String, String>(reader3, reader1Empty, reader2, reader1Empty)
|
def reader = new CompoundReader<String, String>(reader3, reader2)
|
||||||
when:
|
when:
|
||||||
def act = reader.read("test")
|
def act = reader.read("test")
|
||||||
then:
|
then:
|
||||||
StepVerifier.create(act)
|
StepVerifier.create(act)
|
||||||
.expectNext("test-2")
|
.expectNext("test-3")
|
||||||
|
.expectComplete()
|
||||||
|
.verify(Duration.ofSeconds(1))
|
||||||
|
}
|
||||||
|
|
||||||
|
def "Ignore empty"() {
|
||||||
|
setup:
|
||||||
|
def reader = new CompoundReader<String, String>(reader1Empty, reader3, reader2, reader1Empty)
|
||||||
|
when:
|
||||||
|
def act = reader.read("test")
|
||||||
|
then:
|
||||||
|
StepVerifier.create(act)
|
||||||
|
.expectNext("test-3")
|
||||||
.expectComplete()
|
.expectComplete()
|
||||||
.verify(Duration.ofSeconds(1))
|
.verify(Duration.ofSeconds(1))
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user