Elastic schedulers (#349)
This commit is contained in:
@@ -37,8 +37,9 @@ fun main(args: Array<String>) {
|
|||||||
|
|
||||||
HeapDumpCreator.init()
|
HeapDumpCreator.init()
|
||||||
|
|
||||||
|
val cores = Runtime.getRuntime().availableProcessors()
|
||||||
val maxMemory: Long = Runtime.getRuntime().maxMemory() / (1024 * 1024).toLong()
|
val maxMemory: Long = Runtime.getRuntime().maxMemory() / (1024 * 1024).toLong()
|
||||||
log.info("Max heap size: {} MB", maxMemory)
|
log.info("Max heap size: {} MB, number of cores: {}", maxMemory, cores)
|
||||||
|
|
||||||
val app = SpringApplication(Starter::class.java)
|
val app = SpringApplication(Starter::class.java)
|
||||||
app.setDefaultProperties(ResourcePropertySource("version.properties").source)
|
app.setDefaultProperties(ResourcePropertySource("version.properties").source)
|
||||||
|
|||||||
@@ -4,25 +4,33 @@ import io.emeraldpay.dshackle.config.MonitoringConfig
|
|||||||
import io.micrometer.core.instrument.Metrics
|
import io.micrometer.core.instrument.Metrics
|
||||||
import io.micrometer.core.instrument.Tag
|
import io.micrometer.core.instrument.Tag
|
||||||
import io.micrometer.core.instrument.binder.jvm.ExecutorServiceMetrics
|
import io.micrometer.core.instrument.binder.jvm.ExecutorServiceMetrics
|
||||||
|
import org.slf4j.LoggerFactory
|
||||||
import org.springframework.context.annotation.Bean
|
import org.springframework.context.annotation.Bean
|
||||||
import org.springframework.context.annotation.Configuration
|
import org.springframework.context.annotation.Configuration
|
||||||
import org.springframework.scheduling.concurrent.CustomizableThreadFactory
|
import org.springframework.scheduling.concurrent.CustomizableThreadFactory
|
||||||
import reactor.core.scheduler.Scheduler
|
import reactor.core.scheduler.Scheduler
|
||||||
import reactor.core.scheduler.Schedulers
|
import reactor.core.scheduler.Schedulers
|
||||||
import java.util.concurrent.Executor
|
|
||||||
import java.util.concurrent.ExecutorService
|
import java.util.concurrent.ExecutorService
|
||||||
import java.util.concurrent.Executors
|
import java.util.concurrent.Executors
|
||||||
|
|
||||||
@Configuration
|
@Configuration
|
||||||
open class SchedulersConfig {
|
open class SchedulersConfig {
|
||||||
@Bean
|
private val log = LoggerFactory.getLogger(SchedulersConfig::class.java)
|
||||||
open fun rpcScheduler(monitoringConfig: MonitoringConfig): Scheduler {
|
private val threadsMultiplier: Int
|
||||||
return makeScheduler("blockchain-rpc-scheduler", 30, monitoringConfig)
|
|
||||||
|
init {
|
||||||
|
val cores = Runtime.getRuntime().availableProcessors()
|
||||||
|
threadsMultiplier = if (cores < 3) {
|
||||||
|
1
|
||||||
|
} else {
|
||||||
|
cores / 2
|
||||||
|
}
|
||||||
|
log.info("Creating schedulers with multiplier: {}...", threadsMultiplier)
|
||||||
}
|
}
|
||||||
|
|
||||||
@Bean
|
@Bean
|
||||||
open fun trackTxScheduler(monitoringConfig: MonitoringConfig): Scheduler {
|
open fun rpcScheduler(monitoringConfig: MonitoringConfig): Scheduler {
|
||||||
return makeScheduler("tracktx-scheduler", 5, monitoringConfig)
|
return makeScheduler("blockchain-rpc-scheduler", 20, monitoringConfig)
|
||||||
}
|
}
|
||||||
|
|
||||||
@Bean
|
@Bean
|
||||||
@@ -55,18 +63,13 @@ open class SchedulersConfig {
|
|||||||
return makeScheduler("head-liveness-scheduler", 4, monitoringConfig)
|
return makeScheduler("head-liveness-scheduler", 4, monitoringConfig)
|
||||||
}
|
}
|
||||||
|
|
||||||
@Bean
|
|
||||||
open fun grpcChannelExecutor(monitoringConfig: MonitoringConfig): Executor {
|
|
||||||
return makePool("grpc-client-channel", 10, monitoringConfig)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Bean
|
@Bean
|
||||||
open fun authScheduler(monitoringConfig: MonitoringConfig): Scheduler {
|
open fun authScheduler(monitoringConfig: MonitoringConfig): Scheduler {
|
||||||
return makeScheduler("auth-scheduler", 4, monitoringConfig)
|
return makeScheduler("auth-scheduler", 4, monitoringConfig)
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun makeScheduler(name: String, size: Int, monitoringConfig: MonitoringConfig): Scheduler {
|
private fun makeScheduler(name: String, size: Int, monitoringConfig: MonitoringConfig): Scheduler {
|
||||||
return Schedulers.fromExecutorService(makePool(name, size, monitoringConfig))
|
return Schedulers.fromExecutorService(makePool(name, size * threadsMultiplier, monitoringConfig))
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun makePool(name: String, size: Int, monitoringConfig: MonitoringConfig): ExecutorService {
|
private fun makePool(name: String, size: Int, monitoringConfig: MonitoringConfig): ExecutorService {
|
||||||
|
|||||||
Reference in New Issue
Block a user