diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index 3ca9bb2c5c..d3cb8b511b 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -81,6 +81,7 @@ * Add more model installation log info for OAP storage initialization. * BanyanDB: Separate the storage configuration to an independent file: `bydb.yaml`. * Bump Armeria to 1.32.0 and some transitive dependencies. +* Skip persisting metrics/record data that have been expired. #### UI diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java index f07049a285..026f872ba6 100755 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java @@ -29,6 +29,7 @@ import org.apache.skywalking.oap.server.core.analysis.meter.MeterEntity; import org.apache.skywalking.oap.server.core.analysis.meter.MeterSystem; import org.apache.skywalking.oap.server.core.analysis.metrics.ApdexMetrics; import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor; +import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor; import org.apache.skywalking.oap.server.core.analysis.worker.TopNStreamProcessor; import org.apache.skywalking.oap.server.core.annotation.AnnotationScan; import org.apache.skywalking.oap.server.core.cache.AsyncProfilerTaskCache; @@ -379,6 +380,7 @@ public class CoreModuleProvider extends ModuleProvider { metricsStreamProcessor.setL1FlushPeriod(moduleConfig.getL1FlushPeriod()); metricsStreamProcessor.setStorageSessionTimeout(moduleConfig.getStorageSessionTimeout()); metricsStreamProcessor.setMetricsDataTTL(moduleConfig.getMetricsDataTTL()); + RecordStreamProcessor.getInstance().setRecordDataTTL(moduleConfig.getRecordDataTTL()); TopNStreamProcessor.getInstance().setTopNWorkerReportCycle(moduleConfig.getTopNReportPeriod()); apdexThresholdConfig = new ApdexThresholdConfig(this); ApdexMetrics.setDICT(apdexThresholdConfig); 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 a636a72f34..349a60d2fd 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 @@ -102,6 +102,9 @@ public class MetricsPersistentWorker extends PersistenceWorker implemen */ private volatile long timeOfLatestStabilitySts = 0; + // Not going to expose this as a configuration, only for testing purpose + private final boolean isTestingTTL = "true".equalsIgnoreCase(System.getenv("TESTING_TTL")); + MetricsPersistentWorker(ModuleDefineHolder moduleDefineHolder, Model model, IMetricsDAO metricsDAO, AbstractWorker nextAlarmWorker, AbstractWorker nextExportWorker, MetricsTransWorker transWorker, boolean supportUpdate, @@ -188,6 +191,11 @@ public class MetricsPersistentWorker extends PersistenceWorker implemen */ @Override public void in(Metrics metrics) { + final var isExpired = metricsDAO.isExpiredCache(model, metrics, System.currentTimeMillis(), metricsDataTTL); + if (isExpired && !isTestingTTL) { + log.debug("Receiving expired metrics: {}, time: {}, ignored", metrics.id(), metrics.getTimeBucket()); + return; + } aggregationCounter.inc(); dataCarrier.produce(metrics); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/RecordStreamProcessor.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/RecordStreamProcessor.java index a962917e23..06d3230e1b 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/RecordStreamProcessor.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/RecordStreamProcessor.java @@ -21,11 +21,14 @@ package org.apache.skywalking.oap.server.core.analysis.worker; import java.lang.reflect.InvocationTargetException; import java.util.HashMap; import java.util.Map; +import java.util.concurrent.TimeUnit; + import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.UnexpectedException; import org.apache.skywalking.oap.server.core.analysis.DownSampling; import org.apache.skywalking.oap.server.core.analysis.Stream; import org.apache.skywalking.oap.server.core.analysis.StreamProcessor; +import org.apache.skywalking.oap.server.core.analysis.TimeBucket; import org.apache.skywalking.oap.server.core.analysis.record.Record; import org.apache.skywalking.oap.server.core.storage.IRecordDAO; import org.apache.skywalking.oap.server.core.storage.StorageBuilderFactory; @@ -38,18 +41,35 @@ import org.apache.skywalking.oap.server.core.storage.model.ModelCreator; import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder; import org.apache.skywalking.oap.server.library.module.ModuleDefineHolder; +import lombok.Setter; +import lombok.extern.slf4j.Slf4j; + +@Slf4j public class RecordStreamProcessor implements StreamProcessor { private final static RecordStreamProcessor PROCESSOR = new RecordStreamProcessor(); private Map, RecordPersistentWorker> workers = new HashMap<>(); + @Setter + private int recordDataTTL; + + // Not going to expose this as a configuration, only for testing purpose + private final boolean isTestingTTL = "true".equalsIgnoreCase(System.getenv("TESTING_TTL")); + public static RecordStreamProcessor getInstance() { return PROCESSOR; } @Override public void in(Record record) { + final var now = System.currentTimeMillis(); + final var recordTimestamp = TimeBucket.getTimestamp(record.getTimeBucket(), DownSampling.Minute); + final var isExpired = now - recordTimestamp > TimeUnit.DAYS.toMicros(recordDataTTL); + if (isExpired && !isTestingTTL) { + log.debug("Receiving expired record: {}, time: {}, ignored", record.id(), record.getTimeBucket()); + return; + } RecordPersistentWorker worker = workers.get(record.getClass()); if (worker != null) { worker.in(record); diff --git a/test/e2e-v2/cases/cluster/zk/es/docker-compose.yml b/test/e2e-v2/cases/cluster/zk/es/docker-compose.yml index 55a37604c1..5d8e4e9f03 100644 --- a/test/e2e-v2/cases/cluster/zk/es/docker-compose.yml +++ b/test/e2e-v2/cases/cluster/zk/es/docker-compose.yml @@ -59,6 +59,7 @@ services: SW_CLUSTER_INTERNAL_COM_PORT: 11800 SW_CORE_DATA_KEEPER_EXECUTE_PERIOD: 100 SW_TELEMETRY: prometheus + TESTING_TTL: "true" ports: - 1234 depends_on: diff --git a/test/e2e-v2/cases/ttl/es/docker-compose.yml b/test/e2e-v2/cases/ttl/es/docker-compose.yml index 4bf07f62c2..eb23bb09fb 100644 --- a/test/e2e-v2/cases/ttl/es/docker-compose.yml +++ b/test/e2e-v2/cases/ttl/es/docker-compose.yml @@ -43,6 +43,7 @@ services: SW_STORAGE_ES_FLUSH_INTERVAL: 1 SW_CORE_PERSISTENT_PERIOD: 1 SW_CORE_METRICS_DATA_TTL: 7 + TESTING_TTL: "true" depends_on: es: condition: service_healthy