diff --git a/docs/en/setup/backend/telemetry/mesh-mode-grafana.json b/docs/en/setup/backend/telemetry/mesh-mode-grafana.json index 03a2e2650..c1bcb3d61 100644 --- a/docs/en/setup/backend/telemetry/mesh-mode-grafana.json +++ b/docs/en/setup/backend/telemetry/mesh-mode-grafana.json @@ -1074,6 +1074,77 @@ "align": false, "alignLevel": null } + }, + { + "cards": { + "cardPadding": null, + "cardRound": null + }, + "color": { + "cardColor": "#aea2e0", + "colorScale": "sqrt", + "colorScheme": "interpolateOranges", + "exponent": 0.5, + "mode": "opacity" + }, + "dataFormat": "tsbuckets", + "gridPos": { + "h": 6, + "w": 6, + "x": 0, + "y": 27 + }, + "heatmap": {}, + "highlightCards": true, + "id": 24, + "legend": { + "show": false + }, + "links": [], + "repeat": "module", + "repeatDirection": "h", + "scopedVars": { + "module": { + "selected": false, + "text": "endpoint_inventory", + "value": "endpoint_inventory" + } + }, + "targets": [ + { + "expr": "sum(rate(register_persistent_worker_latency_bucket{module=\"[[module]]\"}[10m])) by (le)", + "format": "heatmap", + "hide": false, + "instant": false, + "interval": "15s", + "intervalFactor": 1, + "legendFormat": "{{le}}", + "refId": "A" + } + ], + "title": "register worker latency-$module", + "tooltip": { + "show": true, + "showHistogram": false + }, + "type": "heatmap", + "xAxis": { + "show": true + }, + "xBucketNumber": null, + "xBucketSize": null, + "yAxis": { + "decimals": null, + "format": "s", + "logBase": 1, + "max": null, + "min": null, + "show": true, + "splitFactor": null + }, + "yBucketBound": "auto", + "yBucketNumber": null, + "yBucketSize": null } ], "refresh": false, @@ -1081,7 +1152,34 @@ "style": "dark", "tags": [], "templating": { - "list": [] + "list": [ + { + "allValue": null, + "current": { + "tags": [], + "text": "All", + "value": [ + "$__all" + ] + }, + "datasource": "Prometheus", + "hide": 0, + "includeAll": true, + "label": "module", + "multi": true, + "name": "module", + "options": [], + "query": "label_values(register_persistent_worker_latency_bucket,module)", + "refresh": 1, + "regex": "", + "sort": 0, + "tagValuesQuery": "", + "tags": [], + "tagsQuery": "", + "type": "query", + "useTags": false + } + ] }, "time": { "from": "now-30m", diff --git a/docs/en/setup/backend/telemetry/trace-mode-grafana.json b/docs/en/setup/backend/telemetry/trace-mode-grafana.json index 2bf795946..9cce9b583 100644 --- a/docs/en/setup/backend/telemetry/trace-mode-grafana.json +++ b/docs/en/setup/backend/telemetry/trace-mode-grafana.json @@ -1158,6 +1158,77 @@ "align": false, "alignLevel": null } + }, + { + "cards": { + "cardPadding": null, + "cardRound": null + }, + "color": { + "cardColor": "#aea2e0", + "colorScale": "sqrt", + "colorScheme": "interpolateOranges", + "exponent": 0.5, + "mode": "opacity" + }, + "dataFormat": "tsbuckets", + "gridPos": { + "h": 6, + "w": 6, + "x": 0, + "y": 27 + }, + "heatmap": {}, + "highlightCards": true, + "id": 24, + "legend": { + "show": false + }, + "links": [], + "repeat": "module", + "repeatDirection": "h", + "scopedVars": { + "module": { + "selected": false, + "text": "endpoint_inventory", + "value": "endpoint_inventory" + } + }, + "targets": [ + { + "expr": "sum(rate(register_persistent_worker_latency_bucket{module=\"[[module]]\"}[10m])) by (le)", + "format": "heatmap", + "hide": false, + "instant": false, + "interval": "15s", + "intervalFactor": 1, + "legendFormat": "{{le}}", + "refId": "A" + } + ], + "title": "register worker latency-$module", + "tooltip": { + "show": true, + "showHistogram": false + }, + "type": "heatmap", + "xAxis": { + "show": true + }, + "xBucketNumber": null, + "xBucketSize": null, + "yAxis": { + "decimals": null, + "format": "s", + "logBase": 1, + "max": null, + "min": null, + "show": true, + "splitFactor": null + }, + "yBucketBound": "auto", + "yBucketNumber": null, + "yBucketSize": null } ], "refresh": false, @@ -1165,7 +1236,34 @@ "style": "dark", "tags": [], "templating": { - "list": [] + "list": [ + { + "allValue": null, + "current": { + "tags": [], + "text": "All", + "value": [ + "$__all" + ] + }, + "datasource": "Prometheus", + "hide": 0, + "includeAll": true, + "label": "module", + "multi": true, + "name": "module", + "options": [], + "query": "label_values(register_persistent_worker_latency_bucket,module)", + "refresh": 1, + "regex": "", + "sort": 0, + "tagValuesQuery": "", + "tags": [], + "tagsQuery": "", + "type": "query", + "useTags": false + } + ] }, "time": { "from": "now-5m", diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java index 3a481e948..d58adc2c3 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java @@ -19,6 +19,7 @@ package org.apache.skywalking.oap.server.core.register.worker; import java.util.*; + import org.apache.skywalking.apm.commons.datacarrier.DataCarrier; import org.apache.skywalking.apm.commons.datacarrier.consumer.*; import org.apache.skywalking.oap.server.core.*; @@ -27,6 +28,10 @@ import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine; import org.apache.skywalking.oap.server.core.storage.*; import org.apache.skywalking.oap.server.core.worker.AbstractWorker; import org.apache.skywalking.oap.server.library.module.ModuleDefineHolder; +import org.apache.skywalking.oap.server.telemetry.TelemetryModule; +import org.apache.skywalking.oap.server.telemetry.api.HistogramMetrics; +import org.apache.skywalking.oap.server.telemetry.api.MetricsCreator; +import org.apache.skywalking.oap.server.telemetry.api.MetricsTag; import org.slf4j.*; /** @@ -42,9 +47,10 @@ public class RegisterPersistentWorker extends AbstractWorker { private final IRegisterLockDAO registerLockDAO; private final IRegisterDAO registerDAO; private final DataCarrier dataCarrier; + private final HistogramMetrics workerLatencyHistogram; RegisterPersistentWorker(ModuleDefineHolder moduleDefineHolder, String modelName, - IRegisterDAO registerDAO, int scopeId) { + IRegisterDAO registerDAO, int scopeId) { super(moduleDefineHolder); this.modelName = modelName; this.sources = new HashMap<>(); @@ -52,6 +58,10 @@ public class RegisterPersistentWorker extends AbstractWorker { this.registerLockDAO = moduleDefineHolder.find(StorageModule.NAME).provider().getService(IRegisterLockDAO.class); this.scopeId = scopeId; this.dataCarrier = new DataCarrier<>("MetricsPersistentWorker." + modelName, 1, 1000); + MetricsCreator metricsCreator = moduleDefineHolder.find(TelemetryModule.NAME).provider().getService(MetricsCreator.class); + + workerLatencyHistogram = metricsCreator.createHistogramMetric("register_persistent_worker_latency", "The process latency of register persistent worker", + new MetricsTag.Keys("module"), new MetricsTag.Values(modelName)); String name = "REGISTER_L2"; int size = BulkConsumePool.Creator.recommendMaxSize() / 8; @@ -80,39 +90,41 @@ public class RegisterPersistentWorker extends AbstractWorker { sources.get(registerSource).combine(registerSource); } - if (sources.size() > 1000 || registerSource.isEndOfBatch()) { - sources.values().forEach(source -> { - try { - RegisterSource dbSource = registerDAO.get(modelName, source.id()); - if (Objects.nonNull(dbSource)) { - if (dbSource.combine(source)) { - registerDAO.forceUpdate(modelName, dbSource); - } - } else { - int sequence; - if ((sequence = registerLockDAO.getId(scopeId, source)) != Const.NONE) { - try { - dbSource = registerDAO.get(modelName, source.id()); - if (Objects.nonNull(dbSource)) { - if (dbSource.combine(source)) { - registerDAO.forceUpdate(modelName, dbSource); - } - } else { - source.setSequence(sequence); - registerDAO.forceInsert(modelName, source); - } - } catch (Throwable t) { - logger.error(t.getMessage(), t); + try (HistogramMetrics.Timer timer = workerLatencyHistogram.createTimer()) { + if (sources.size() > 1000 || registerSource.isEndOfBatch()) { + sources.values().forEach(source -> { + try { + RegisterSource dbSource = registerDAO.get(modelName, source.id()); + if (Objects.nonNull(dbSource)) { + if (dbSource.combine(source)) { + registerDAO.forceUpdate(modelName, dbSource); } } else { - logger.info("{} inventory register try lock and increment sequence failure.", DefaultScopeDefine.nameOf(scopeId)); + int sequence; + if ((sequence = registerLockDAO.getId(scopeId, source)) != Const.NONE) { + try { + dbSource = registerDAO.get(modelName, source.id()); + if (Objects.nonNull(dbSource)) { + if (dbSource.combine(source)) { + registerDAO.forceUpdate(modelName, dbSource); + } + } else { + source.setSequence(sequence); + registerDAO.forceInsert(modelName, source); + } + } catch (Throwable t) { + logger.error(t.getMessage(), t); + } + } else { + logger.info("{} inventory register try lock and increment sequence failure.", DefaultScopeDefine.nameOf(scopeId)); + } } + } catch (Throwable t) { + logger.error(t.getMessage(), t); } - } catch (Throwable t) { - logger.error(t.getMessage(), t); - } - }); - sources.clear(); + }); + sources.clear(); + } } } diff --git a/oap-server/server-telemetry/telemetry-api/src/main/java/org/apache/skywalking/oap/server/telemetry/api/HistogramMetrics.java b/oap-server/server-telemetry/telemetry-api/src/main/java/org/apache/skywalking/oap/server/telemetry/api/HistogramMetrics.java index 41e73da20..dce0af347 100644 --- a/oap-server/server-telemetry/telemetry-api/src/main/java/org/apache/skywalking/oap/server/telemetry/api/HistogramMetrics.java +++ b/oap-server/server-telemetry/telemetry-api/src/main/java/org/apache/skywalking/oap/server/telemetry/api/HistogramMetrics.java @@ -55,7 +55,7 @@ public abstract class HistogramMetrics { } @Override - public void close() throws IOException { + public void close() { finish(); } }