Merge pull request #112 from p2p-org/fix-extra-cpu-usage-2

Fix extra cpu usage
- added fixed channel executor thread pool
- lag observer should be created only in case there are more than 1 upstreams exists
- fixed small typos in thread pools names
This commit is contained in:
a10zn8
2023-01-17 20:21:37 +04:00
committed by GitHub
6 changed files with 49 additions and 24 deletions

View File

@@ -66,7 +66,7 @@ open class GrpcServer(
serverBuilder.addService(it)
}
val pool = Executors.newFixedThreadPool(20, CustomizableThreadFactory("fixed-grpc-%d"))
val pool = Executors.newFixedThreadPool(20, CustomizableThreadFactory("fixed-grpc-"))
serverBuilder.executor(
if (mainConfig.monitoring.enableExtended)

View File

@@ -8,6 +8,8 @@ import org.springframework.context.annotation.Configuration
import org.springframework.scheduling.concurrent.CustomizableThreadFactory
import reactor.core.scheduler.Scheduler
import reactor.core.scheduler.Schedulers
import java.util.concurrent.Executor
import java.util.concurrent.ExecutorService
import java.util.concurrent.Executors
@Configuration
@@ -27,19 +29,26 @@ open class SchedulersConfig {
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"))
@Bean
open fun grpcChannelExecutor(monitoringConfig: MonitoringConfig): Executor {
return makePool("grpc-client-channel", "grpc_client_channel", 10, monitoringConfig)
}
return Schedulers.fromExecutorService(
if (monitoringConfig.enableExtended)
ExecutorServiceMetrics.monitor(
Metrics.globalRegistry,
pool,
name,
prefix
)
else
pool
)
private fun makeScheduler(name: String, prefix: String, size: Int, monitoringConfig: MonitoringConfig): Scheduler {
return Schedulers.fromExecutorService(makePool(name, prefix, size, monitoringConfig))
}
private fun makePool(name: String, prefix: String, size: Int, monitoringConfig: MonitoringConfig): ExecutorService {
val pool = Executors.newFixedThreadPool(size, CustomizableThreadFactory("$name-"))
return if (monitoringConfig.enableExtended)
ExecutorServiceMetrics.monitor(
Metrics.globalRegistry,
pool,
name,
prefix
)
else
pool
}
}

View File

@@ -42,6 +42,7 @@ import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.forkchoice.NoChoiceWithPriorityForkChoice
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstreams
import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Qualifier
import org.springframework.boot.ApplicationArguments
import org.springframework.boot.ApplicationRunner
import org.springframework.context.ApplicationEventPublisher
@@ -50,6 +51,7 @@ import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import reactor.core.scheduler.Schedulers
import java.net.URI
import java.util.concurrent.Executor
import java.util.concurrent.atomic.AtomicInteger
import java.util.function.Function
import kotlin.math.abs
@@ -59,7 +61,9 @@ open class ConfiguredUpstreams(
private val fileResolver: FileResolver,
private val config: UpstreamsConfig,
private val callTargets: CallTargetsHolder,
private val eventPublisher: ApplicationEventPublisher
private val eventPublisher: ApplicationEventPublisher,
@Qualifier("grpcChannelExecutor")
private val channelExecutor: Executor
) : ApplicationRunner {
private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java)
@@ -308,7 +312,8 @@ open class ConfiguredUpstreams(
endpoint.auth,
fileResolver,
endpoint.upstreamRating,
config.labels
config.labels,
channelExecutor
).apply {
timeout = options.timeout
}

View File

@@ -198,8 +198,9 @@ abstract class Multistream(
}
lagObserver?.stop()
lagObserver = null
if (upstreams.isNotEmpty()) {
lagObserver = makeLagObserver()
when {
upstreams.size == 1 -> upstreams[0].setLag(0)
upstreams.size > 1 -> lagObserver = makeLagObserver()
}
}
}

View File

@@ -45,6 +45,7 @@ import reactor.core.Disposable
import reactor.core.publisher.Flux
import java.io.IOException
import java.time.Duration
import java.util.concurrent.Executor
import java.util.concurrent.atomic.AtomicReference
import java.util.concurrent.locks.ReentrantLock
import kotlin.concurrent.withLock
@@ -58,7 +59,8 @@ class GrpcUpstreams(
private val auth: AuthConfig.ClientTlsAuth? = null,
private val fileResolver: FileResolver,
private val nodeRating: Int,
private val labels: UpstreamsConfig.Labels
private val labels: UpstreamsConfig.Labels,
private val executor: Executor
) {
private val log = LoggerFactory.getLogger(GrpcUpstreams::class.java)
@@ -73,6 +75,7 @@ class GrpcUpstreams(
// some messages are very large. many of them in megabytes, some even in gigabytes (ex. ETH Traces)
.maxInboundMessageSize(Defaults.maxMessageSize)
.enableRetry()
.executor(executor)
.maxRetryAttempts(3)
if (auth != null && StringUtils.isNotEmpty(auth.ca)) {
chanelBuilder

View File

@@ -11,6 +11,8 @@ import io.emeraldpay.dshackle.Chain
import org.springframework.context.ApplicationEventPublisher
import spock.lang.Specification
import java.util.concurrent.Executors
class ConfiguredUpstreamsSpec extends Specification {
def "Applied quorum to extra methods"() {
@@ -20,7 +22,8 @@ class ConfiguredUpstreamsSpec extends Specification {
Stub(FileResolver),
Stub(UpstreamsConfig),
callTargetsHolder,
Mock(ApplicationEventPublisher)
Mock(ApplicationEventPublisher),
Executors.newFixedThreadPool(1)
)
def methods = new UpstreamsConfig.Methods(
[
@@ -45,7 +48,8 @@ class ConfiguredUpstreamsSpec extends Specification {
Stub(FileResolver),
Stub(UpstreamsConfig),
callTargetsHolder,
Mock(ApplicationEventPublisher)
Mock(ApplicationEventPublisher),
Executors.newFixedThreadPool(1)
)
def methods = new UpstreamsConfig.Methods(
[
@@ -69,7 +73,8 @@ class ConfiguredUpstreamsSpec extends Specification {
Stub(FileResolver),
Stub(UpstreamsConfig),
callTargetsHolder,
Mock(ApplicationEventPublisher)
Mock(ApplicationEventPublisher),
Executors.newFixedThreadPool(1)
)
expect:
configurer.getHash(node, src) == expected
@@ -88,7 +93,8 @@ class ConfiguredUpstreamsSpec extends Specification {
Stub(FileResolver),
Stub(UpstreamsConfig),
callTargetsHolder,
Mock(ApplicationEventPublisher)
Mock(ApplicationEventPublisher),
Executors.newFixedThreadPool(1)
)
when:
def h1 = configurer.getHash(null, "hohoho")
@@ -112,7 +118,8 @@ class ConfiguredUpstreamsSpec extends Specification {
Stub(FileResolver),
Stub(UpstreamsConfig),
callTargetsHolder,
Mock(ApplicationEventPublisher)
Mock(ApplicationEventPublisher),
Executors.newFixedThreadPool(1)
)
def methodsGroup = new UpstreamsConfig.MethodGroups(
["filter"] as Set,