Small improvements after dynamic head introduction
- add injection of schedulers for dynamic heads - more logging - replace onBackpressureBuffer to directBestEffort multicast sink
This commit is contained in:
@@ -8,7 +8,7 @@ import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
class DynamicMergeFlux<K : Any, T>(private val scheduler: Scheduler) {
|
||||
|
||||
private val merge = Sinks.many().multicast().onBackpressureBuffer<T>()
|
||||
private val merge = Sinks.many().multicast().directBestEffort<T>()
|
||||
private val sources = ConcurrentHashMap<K, Disposable>()
|
||||
|
||||
fun add(flux: Flux<T>, id: K) {
|
||||
@@ -30,4 +30,6 @@ class DynamicMergeFlux<K : Any, T>(private val scheduler: Scheduler) {
|
||||
sources.forEach { (_, d) -> d.dispose() }
|
||||
merge.emitComplete { _, res -> res == Sinks.EmitResult.FAIL_NON_SERIALIZED }
|
||||
}
|
||||
|
||||
fun getKeys(): List<K> = sources.keys().toList()
|
||||
}
|
||||
|
||||
@@ -8,23 +8,27 @@ import io.emeraldpay.dshackle.upstream.Multistream
|
||||
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinMultistream
|
||||
import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream
|
||||
import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream
|
||||
import org.springframework.beans.factory.annotation.Qualifier
|
||||
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory
|
||||
import org.springframework.context.annotation.Bean
|
||||
import org.springframework.context.annotation.Configuration
|
||||
import reactor.core.scheduler.Scheduler
|
||||
|
||||
@Configuration
|
||||
open class MultistreamsConfig(val beanFactory: ConfigurableListableBeanFactory) {
|
||||
@Bean
|
||||
open fun allMultistreams(
|
||||
cachesFactory: CachesFactory,
|
||||
callTargetsHolder: CallTargetsHolder
|
||||
callTargetsHolder: CallTargetsHolder,
|
||||
@Qualifier("headMergedScheduler")
|
||||
headScheduler: Scheduler
|
||||
): List<Multistream> {
|
||||
return Chain.values()
|
||||
.filterNot { it == Chain.UNSPECIFIED }
|
||||
.mapNotNull { chain ->
|
||||
when (BlockchainType.from(chain)) {
|
||||
BlockchainType.EVM_POS -> ethereumPosMultistream(chain, cachesFactory)
|
||||
BlockchainType.EVM_POW -> ethereumMultistream(chain, cachesFactory)
|
||||
BlockchainType.EVM_POS -> ethereumPosMultistream(chain, cachesFactory, headScheduler)
|
||||
BlockchainType.EVM_POW -> ethereumMultistream(chain, cachesFactory, headScheduler)
|
||||
BlockchainType.BITCOIN -> bitcoinMultistream(chain, cachesFactory)
|
||||
else -> null
|
||||
}
|
||||
@@ -33,27 +37,31 @@ open class MultistreamsConfig(val beanFactory: ConfigurableListableBeanFactory)
|
||||
|
||||
private fun ethereumMultistream(
|
||||
chain: Chain,
|
||||
cachesFactory: CachesFactory
|
||||
cachesFactory: CachesFactory,
|
||||
headScheduler: Scheduler
|
||||
): EthereumMultistream {
|
||||
val name = "multi-ethereum-$chain"
|
||||
|
||||
return EthereumMultistream(
|
||||
chain,
|
||||
ArrayList(),
|
||||
cachesFactory.getCaches(chain)
|
||||
cachesFactory.getCaches(chain),
|
||||
headScheduler
|
||||
).also { register(it, name) }
|
||||
}
|
||||
|
||||
open fun ethereumPosMultistream(
|
||||
chain: Chain,
|
||||
cachesFactory: CachesFactory
|
||||
cachesFactory: CachesFactory,
|
||||
headScheduler: Scheduler
|
||||
): EthereumPosMultiStream {
|
||||
val name = "multi-ethereum-pos-$chain"
|
||||
|
||||
return EthereumPosMultiStream(
|
||||
chain,
|
||||
ArrayList(),
|
||||
cachesFactory.getCaches(chain)
|
||||
cachesFactory.getCaches(chain),
|
||||
headScheduler
|
||||
).also { register(it, name) }
|
||||
}
|
||||
|
||||
|
||||
@@ -22,6 +22,11 @@ open class SchedulersConfig {
|
||||
return makeScheduler("tracktx-scheduler", "tracktx", 5, monitoringConfig)
|
||||
}
|
||||
|
||||
@Bean
|
||||
open fun headMergedScheduler(monitoringConfig: MonitoringConfig): Scheduler {
|
||||
return makeScheduler("head-scheduler", "head_merge", 5, monitoringConfig)
|
||||
}
|
||||
|
||||
private fun makeScheduler(name: String, prefix: String, size: Int, monitoringConfig: MonitoringConfig): Scheduler {
|
||||
val pool = Executors.newFixedThreadPool(size, CustomizableThreadFactory("$name-%d"))
|
||||
|
||||
|
||||
@@ -40,10 +40,12 @@ open class DynamicMergedHead(
|
||||
}
|
||||
|
||||
fun addHead(upstream: Upstream) {
|
||||
log.debug("adding upstream head of [${upstream.getId()}] to dynamic head of [$label]. Current heads ${dynamicFlux.getKeys()}")
|
||||
dynamicFlux.add(upstream.getHead().getFlux(), upstream.getId())
|
||||
}
|
||||
|
||||
fun removeHead(id: String) {
|
||||
log.debug("removing upstream head of [$id] from dynamic head $label. Current heads ${dynamicFlux.getKeys()}")
|
||||
dynamicFlux.remove(id)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,13 +32,14 @@ import org.slf4j.LoggerFactory
|
||||
import org.springframework.util.ConcurrentReferenceHashMap
|
||||
import reactor.core.publisher.Flux
|
||||
import reactor.core.publisher.Mono
|
||||
import reactor.core.scheduler.Schedulers
|
||||
import reactor.core.scheduler.Scheduler
|
||||
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
open class EthereumMultistream(
|
||||
chain: Chain,
|
||||
val upstreams: MutableList<EthereumUpstream>,
|
||||
caches: Caches
|
||||
caches: Caches,
|
||||
headScheduler: Scheduler
|
||||
) : Multistream(chain, upstreams as MutableList<Upstream>, caches), EthereumLikeMultistream {
|
||||
|
||||
companion object {
|
||||
@@ -48,7 +49,7 @@ open class EthereumMultistream(
|
||||
private var head: DynamicMergedHead = DynamicMergedHead(
|
||||
PriorityForkChoice(),
|
||||
"ETH Multistream",
|
||||
Schedulers.boundedElastic()
|
||||
headScheduler
|
||||
)
|
||||
|
||||
private val filteredHeads: MutableMap<String, Head> =
|
||||
|
||||
@@ -31,13 +31,14 @@ import org.slf4j.LoggerFactory
|
||||
import org.springframework.util.ConcurrentReferenceHashMap
|
||||
import reactor.core.publisher.Flux
|
||||
import reactor.core.publisher.Mono
|
||||
import reactor.core.scheduler.Schedulers
|
||||
import reactor.core.scheduler.Scheduler
|
||||
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
open class EthereumPosMultiStream(
|
||||
chain: Chain,
|
||||
val upstreams: MutableList<EthereumPosUpstream>,
|
||||
caches: Caches
|
||||
caches: Caches,
|
||||
headScheduler: Scheduler
|
||||
) : Multistream(chain, upstreams as MutableList<Upstream>, caches), EthereumLikeMultistream {
|
||||
|
||||
companion object {
|
||||
@@ -47,7 +48,7 @@ open class EthereumPosMultiStream(
|
||||
private var head: DynamicMergedHead = DynamicMergedHead(
|
||||
PriorityForkChoice(),
|
||||
"ETH Pos Multistream",
|
||||
Schedulers.boundedElastic()
|
||||
headScheduler
|
||||
)
|
||||
|
||||
private val reader: EthereumCachingReader = EthereumCachingReader(this, this.caches, getMethodsFactory())
|
||||
|
||||
Reference in New Issue
Block a user