added metrics for native call inner requests
This commit is contained in:
@@ -32,9 +32,7 @@ import org.slf4j.LoggerFactory
|
|||||||
import reactor.netty.http.server.HttpServer
|
import reactor.netty.http.server.HttpServer
|
||||||
import reactor.netty.http.server.HttpServerRoutes
|
import reactor.netty.http.server.HttpServerRoutes
|
||||||
import java.util.EnumMap
|
import java.util.EnumMap
|
||||||
import java.util.concurrent.locks.ReentrantReadWriteLock
|
import java.util.concurrent.ConcurrentHashMap
|
||||||
import kotlin.concurrent.read
|
|
||||||
import kotlin.concurrent.write
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* HTTP Proxy Server
|
* HTTP Proxy Server
|
||||||
@@ -141,32 +139,13 @@ class ProxyServer(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* Monitoring that has separate metrics per RPC method.
|
|
||||||
* Slightly slower to use than StandardRequestMetrics
|
|
||||||
*/
|
|
||||||
class ExtendedRequestMetrics : RequestMetricsFactory {
|
class ExtendedRequestMetrics : RequestMetricsFactory {
|
||||||
private val current = HashMap<Ref, RequestMetrics>()
|
private val current = ConcurrentHashMap<Ref, RequestMetrics>()
|
||||||
private val lock = ReentrantReadWriteLock()
|
|
||||||
|
|
||||||
override fun get(chain: Chain, method: String): RequestMetrics {
|
override fun get(chain: Chain, method: String): RequestMetrics =
|
||||||
val ref = Ref(chain, method)
|
current.computeIfAbsent(Ref(chain, method)) {
|
||||||
lock.read {
|
RequestMetricsWithMethod(it.chain, it.method)
|
||||||
val existing = current[ref]
|
|
||||||
if (existing != null) {
|
|
||||||
return existing
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
lock.write {
|
|
||||||
val existing = current[ref]
|
|
||||||
if (existing != null) {
|
|
||||||
return existing
|
|
||||||
}
|
|
||||||
val created = RequestMetricsWithMethod(chain, method)
|
|
||||||
current[ref] = created
|
|
||||||
return created
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
data class Ref(val chain: Chain, val method: String)
|
data class Ref(val chain: Chain, val method: String)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -32,6 +32,7 @@ import org.springframework.stereotype.Service
|
|||||||
import reactor.core.publisher.Flux
|
import reactor.core.publisher.Flux
|
||||||
import reactor.core.publisher.Mono
|
import reactor.core.publisher.Mono
|
||||||
import java.util.Locale
|
import java.util.Locale
|
||||||
|
import java.util.concurrent.ConcurrentHashMap
|
||||||
import java.util.concurrent.TimeUnit
|
import java.util.concurrent.TimeUnit
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
@@ -65,18 +66,28 @@ class BlockchainRpc(
|
|||||||
override fun nativeCall(request: Mono<BlockchainOuterClass.NativeCallRequest>): Flux<BlockchainOuterClass.NativeCallReplyItem> {
|
override fun nativeCall(request: Mono<BlockchainOuterClass.NativeCallRequest>): Flux<BlockchainOuterClass.NativeCallReplyItem> {
|
||||||
var startTime = 0L
|
var startTime = 0L
|
||||||
var metrics: RequestMetrics? = null
|
var metrics: RequestMetrics? = null
|
||||||
|
val idsMap = mutableMapOf<Int, String>()
|
||||||
return nativeCall.nativeCall(
|
return nativeCall.nativeCall(
|
||||||
request
|
request
|
||||||
.doOnNext {
|
.doOnNext { req ->
|
||||||
metrics = chainMetrics.get(it.chain)
|
metrics = chainMetrics.get(req.chain)
|
||||||
metrics!!.nativeCallMetric.increment()
|
metrics?.let { m ->
|
||||||
|
m.nativeCallMetric.increment()
|
||||||
|
req.itemsList.forEach { item ->
|
||||||
|
idsMap[item.id] = item.method
|
||||||
|
m.getNativeItemMetrics(item.method).nativeItemRequest.increment()
|
||||||
|
}
|
||||||
|
}
|
||||||
startTime = System.currentTimeMillis()
|
startTime = System.currentTimeMillis()
|
||||||
}
|
}
|
||||||
).doOnNext { reply ->
|
).doOnNext { reply ->
|
||||||
metrics?.let { m ->
|
metrics?.getNativeItemMetrics(idsMap[reply.id] ?: "unknown")?.let { itemMetrics ->
|
||||||
m.nativeCallRespMetric?.record(System.currentTimeMillis() - startTime, TimeUnit.MILLISECONDS)
|
itemMetrics.nativeItemResponse.record(
|
||||||
|
System.currentTimeMillis() - startTime,
|
||||||
|
TimeUnit.MILLISECONDS
|
||||||
|
)
|
||||||
if (!reply.succeed) {
|
if (!reply.succeed) {
|
||||||
m.nativeCallErrRespMetric.increment()
|
itemMetrics.nativeItemResponseErr.increment()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}.doOnError { failMetric.increment() }
|
}.doOnError { failMetric.increment() }
|
||||||
@@ -204,20 +215,11 @@ class BlockchainRpc(
|
|||||||
.doOnError { failMetric.increment() }
|
.doOnError { failMetric.increment() }
|
||||||
}
|
}
|
||||||
|
|
||||||
class RequestMetrics(chain: Chain) {
|
class RequestMetrics(val chain: Chain) {
|
||||||
val nativeCallMetric = Counter.builder("request.grpc.request")
|
val nativeCallMetric = Counter.builder("request.grpc.request")
|
||||||
.tag("type", "nativeCall")
|
.tag("type", "nativeCall")
|
||||||
.tag("chain", chain.chainCode)
|
.tag("chain", chain.chainCode)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
val nativeCallRespMetric = Timer.builder("request.grpc.response")
|
|
||||||
.tag("type", "nativeCall")
|
|
||||||
.tag("chain", chain.chainCode)
|
|
||||||
.publishPercentileHistogram()
|
|
||||||
.register(Metrics.globalRegistry)
|
|
||||||
val nativeCallErrRespMetric = Counter.builder("request.grpc.response.err")
|
|
||||||
.tag("type", "nativeCall")
|
|
||||||
.tag("chain", chain.chainCode)
|
|
||||||
.register(Metrics.globalRegistry)
|
|
||||||
val nativeSubscribeMetric = Counter.builder("request.grpc.request")
|
val nativeSubscribeMetric = Counter.builder("request.grpc.request")
|
||||||
.tag("type", "nativeSubscribe")
|
.tag("type", "nativeSubscribe")
|
||||||
.tag("chain", chain.chainCode)
|
.tag("chain", chain.chainCode)
|
||||||
@@ -264,5 +266,26 @@ class BlockchainRpc(
|
|||||||
.tag("chain", chain.chainCode)
|
.tag("chain", chain.chainCode)
|
||||||
.publishPercentileHistogram()
|
.publishPercentileHistogram()
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
|
|
||||||
|
private val nativeItemMetrics = ConcurrentHashMap<String, NativeRequestItemsMetrics>()
|
||||||
|
|
||||||
|
fun getNativeItemMetrics(method: String) =
|
||||||
|
nativeItemMetrics.computeIfAbsent(method) { NativeRequestItemsMetrics(chain, it) }
|
||||||
|
}
|
||||||
|
|
||||||
|
class NativeRequestItemsMetrics(chain: Chain, method: String) {
|
||||||
|
val nativeItemRequest = Counter.builder("request.grpc.native.request")
|
||||||
|
.tag("chain", chain.chainCode)
|
||||||
|
.tag("method", method)
|
||||||
|
.register(Metrics.globalRegistry)
|
||||||
|
val nativeItemResponse = Timer.builder("request.grpc.native.response")
|
||||||
|
.tag("chain", chain.chainCode)
|
||||||
|
.tag("method", method)
|
||||||
|
.publishPercentileHistogram()
|
||||||
|
.register(Metrics.globalRegistry)
|
||||||
|
val nativeItemResponseErr = Counter.builder("request.grpc.native.request")
|
||||||
|
.tag("chain", chain.chainCode)
|
||||||
|
.tag("method", method)
|
||||||
|
.register(Metrics.globalRegistry)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user