Add method label to upstream.rpc.conn metric (#863)
This commit is contained in:
@@ -34,11 +34,14 @@ class BasicHttpFactory(
|
||||
Tag.of("chain", chain.chainCode),
|
||||
)
|
||||
val metrics = RequestMetrics(
|
||||
Timer.builder("upstream.rpc.conn")
|
||||
.description("Request time through a HTTP JSON RPC connection")
|
||||
.tags(metricsTags)
|
||||
.publishPercentileHistogram()
|
||||
.register(Metrics.globalRegistry),
|
||||
{ method ->
|
||||
Timer.builder("upstream.rpc.conn")
|
||||
.description("Request time through a HTTP JSON RPC connection")
|
||||
.tags(metricsTags)
|
||||
.tag("method", method ?: "unknown")
|
||||
.publishPercentileHistogram()
|
||||
.register(Metrics.globalRegistry)
|
||||
},
|
||||
Counter.builder("upstream.rpc.fail")
|
||||
.description("Number of failures of HTTP JSON RPC requests")
|
||||
.tags(metricsTags)
|
||||
|
||||
@@ -101,7 +101,7 @@ abstract class HttpReader(
|
||||
|
||||
open fun onStop() {
|
||||
if (metrics != null) {
|
||||
Metrics.globalRegistry.remove(metrics.timer)
|
||||
metrics.registeredTimers().forEach { Metrics.globalRegistry.remove(it) }
|
||||
Metrics.globalRegistry.remove(metrics.fails)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,9 +17,20 @@ package io.emeraldpay.dshackle.upstream
|
||||
|
||||
import io.micrometer.core.instrument.Counter
|
||||
import io.micrometer.core.instrument.Timer
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
class RequestMetrics(
|
||||
val timer: Timer,
|
||||
private val timerFactory: (String?) -> Timer,
|
||||
val fails: Counter,
|
||||
val nettyMetricsEnabled: Boolean,
|
||||
)
|
||||
) {
|
||||
private val timerCache = ConcurrentHashMap<String, Timer>()
|
||||
|
||||
constructor(timer: Timer, fails: Counter, nettyMetricsEnabled: Boolean) :
|
||||
this({ _ -> timer }, fails, nettyMetricsEnabled)
|
||||
|
||||
fun timer(method: String? = null): Timer =
|
||||
timerCache.computeIfAbsent(method ?: "") { timerFactory(method) }
|
||||
|
||||
fun registeredTimers(): Collection<Timer> = timerCache.values
|
||||
}
|
||||
|
||||
@@ -404,7 +404,7 @@ open class WsConnectionImpl(
|
||||
return Mono.from(onResponse.asMono()).or(failOnDisconnect)
|
||||
.doOnSubscribe { sendRpc(request) }
|
||||
.take(Defaults.timeout)
|
||||
.doOnNext { requestMetrics?.timer?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) }
|
||||
.doOnNext { requestMetrics?.timer()?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) }
|
||||
.doOnError { requestMetrics?.fails?.increment() }
|
||||
.map { it.copyWithId(ChainResponse.Id.from(originalId)) }
|
||||
.switchIfEmpty(
|
||||
@@ -424,7 +424,7 @@ open class WsConnectionImpl(
|
||||
it.close()
|
||||
Metrics.globalRegistry.remove(it)
|
||||
}
|
||||
requestMetrics?.timer?.let {
|
||||
requestMetrics?.registeredTimers()?.forEach {
|
||||
it.close()
|
||||
Metrics.globalRegistry.remove(it)
|
||||
}
|
||||
|
||||
@@ -59,7 +59,7 @@ class RestHttpReader(
|
||||
.flatMap(this::execute)
|
||||
.doOnNext {
|
||||
if (startTime.isStarted) {
|
||||
metrics?.timer?.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
|
||||
metrics?.timer(key.method)?.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
|
||||
}
|
||||
}
|
||||
.handle { it, sink ->
|
||||
|
||||
@@ -114,7 +114,7 @@ class JsonRpcGrpcClient(
|
||||
}
|
||||
.doOnNext {
|
||||
if (timer.isStarted) {
|
||||
metrics?.timer?.record(timer.getTime(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS)
|
||||
metrics?.timer()?.record(timer.getTime(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -92,7 +92,7 @@ class JsonRpcHttpReader(
|
||||
.flatMap(this@JsonRpcHttpReader::execute)
|
||||
.doOnNext {
|
||||
if (startTime.isStarted) {
|
||||
metrics?.timer?.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
|
||||
metrics?.timer(key.method)?.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
|
||||
}
|
||||
}
|
||||
.transform(asJsonRpcResponse(key))
|
||||
|
||||
Reference in New Issue
Block a user