diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index 1508721700..a0e5585037 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -71,6 +71,7 @@ * Bump up netty to 4.11.118 to fix CVE-2025-24970. * Add `Get Alarm Runtime Status` API. * Add `lock` when query the Alarm metrics window values. +* Add a fail-safe mechanism to prevent traffic metrics inconsistent between in-memory and database server. #### UI diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsPersistentWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsPersistentWorker.java index 822540a598..2f9029fd79 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsPersistentWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsPersistentWorker.java @@ -107,7 +107,7 @@ public class MetricsPersistentWorker extends PersistenceWorker implemen long storageSessionTimeout, int metricsDataTTL, MetricStreamKind kind) { super(moduleDefineHolder, new ReadWriteSafeCache<>(new MergableBufferedData(), new MergableBufferedData())); this.model = model; - this.sessionCache = new MetricsSessionCache(storageSessionTimeout); + this.sessionCache = new MetricsSessionCache(storageSessionTimeout, supportUpdate); this.metricsDAO = metricsDAO; this.nextAlarmWorker = Optional.ofNullable(nextAlarmWorker); this.nextExportWorker = Optional.ofNullable(nextExportWorker); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsSessionCache.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsSessionCache.java index 93b1e421cf..1b8da1e28d 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsSessionCache.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsSessionCache.java @@ -19,16 +19,23 @@ package org.apache.skywalking.oap.server.core.analysis.worker; import java.util.Iterator; +import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import lombok.AccessLevel; import lombok.Setter; +import org.apache.skywalking.oap.server.core.analysis.MetricsExtension; +import org.apache.skywalking.oap.server.core.analysis.manual.service.ServiceTraffic; import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics; +import org.apache.skywalking.oap.server.core.storage.IMetricsDAO; +import org.apache.skywalking.oap.server.core.storage.SessionCacheCallback; +import org.apache.skywalking.oap.server.core.storage.model.Model; import org.apache.skywalking.oap.server.library.client.request.InsertRequest; import org.apache.skywalking.oap.server.library.client.request.UpdateRequest; /** * MetricsSessionCache is a key-value cache to hold hot metric in-memory to reduce payload to pre-read. + * Every instance of MetricsSessionCache maps to an instance of {@link MetricsPersistentWorker}. * * There are two ways to make sure metrics in-cache, * 1. Metrics is read from the Database through {@link MetricsPersistentWorker}.loadFromStorage @@ -45,14 +52,16 @@ public class MetricsSessionCache { private final Map sessionCache; @Setter(AccessLevel.PACKAGE) private long timeoutThreshold; + private final boolean supportUpdate; - public MetricsSessionCache(long timeoutThreshold) { + public MetricsSessionCache(long timeoutThreshold, final boolean supportUpdate) { // Due to the cache would be updated depending on final storage implementation, // the map/cache could be updated concurrently. // Set to ConcurrentHashMap in order to avoid HashMap deadlock. // Since 9.3.0 this.sessionCache = new ConcurrentHashMap<>(100); this.timeoutThreshold = timeoutThreshold; + this.supportUpdate = supportUpdate; } Metrics get(Metrics metrics) { @@ -67,6 +76,24 @@ public class MetricsSessionCache { sessionCache.put(metrics, metrics); } + /** + * This method relies on the response of database flush callback. + * Push the data into the in-memory cache for all metrics except {@link MetricsExtension#supportUpdate()} labeled + * as false. + * Because those data(e.g. {@link ServiceTraffic}) is one-time writing in the whole TTL period, and some + * database(e.g. BanyanDB) has in-memory cache at the server side to improve performance but trade off the 100% + * eventual consistency of writing, which means database server could respond + * {@link SessionCacheCallback#onInsertCompleted()} but ends of writing failure caused by crashing. + * This fail-safe mechanism would require the cache of this kind of metric must be read through + * {@link IMetricsDAO#multiGet(Model, List)} which guaranteed data existence. + */ + public void cacheAfterFlush(Metrics metrics) { + if (supportUpdate) { + put(metrics); + } + // Don't add into the cache. Rely on multiGet from Database. + } + void removeExpired() { Iterator iterator = sessionCache.values().iterator(); long timestamp = System.currentTimeMillis(); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/SessionCacheCallback.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/SessionCacheCallback.java index ceff064e05..aa0149737b 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/SessionCacheCallback.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/SessionCacheCallback.java @@ -40,7 +40,7 @@ public class SessionCacheCallback { if (isFailed) { return; } - sessionCache.put(metrics); + sessionCache.cacheAfterFlush(metrics); } public void onUpdateFailure() {