From 8dc69e3b5d10864989b8224e5537a84c7da8fb28 Mon Sep 17 00:00:00 2001 From: waterWang <672684719@qq.com> Date: Sat, 22 Aug 2026 11:56:05 +0800 Subject: [PATCH] fix: use epoch-based monotonic message ID for BanyanDB measure writes Replace System.nanoTime() with a message ID generator that uses the Unix epoch as its base. The generator is backed by an AtomicLong that guarantees strictly increasing values even under concurrent writes or a wall-clock rollback, making message IDs comparable across OAP hosts and process restarts. For MeasureWrite, the generated value is also set as DataPointValue.version so that BanyanDB does not need to derive it from a host-local source. Fixes #13986 --- .../banyandb/v1/client/AbstractWrite.java | 17 +++ .../banyandb/v1/client/MeasureWrite.java | 8 +- .../banyandb/v1/client/StreamWrite.java | 4 +- .../v1/client/MeasureWriteMessageIdTest.java | 101 ++++++++++++++++++ 4 files changed, 126 insertions(+), 4 deletions(-) create mode 100644 oap-server/server-library/library-banyandb-client/src/test/java/org/apache/skywalking/library/banyandb/v1/client/MeasureWriteMessageIdTest.java diff --git a/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/AbstractWrite.java b/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/AbstractWrite.java index 73cd0f11bdc5..64b577cf64a7 100644 --- a/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/AbstractWrite.java +++ b/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/AbstractWrite.java @@ -19,10 +19,27 @@ package org.apache.skywalking.library.banyandb.v1.client; import java.util.Optional; +import java.util.concurrent.atomic.AtomicLong; import lombok.Getter; import org.apache.skywalking.banyandb.common.v1.BanyandbCommon; public abstract class AbstractWrite
{ + /** + * Generates monotonic message IDs for write requests. + *
+ * The value is based on the Unix epoch (wall clock) so that message IDs are comparable across + * OAP hosts and process restarts, instead of {@link System#nanoTime()} which is only meaningful + * within a single JVM instance. An {@link AtomicLong} keeps the value strictly increasing even + * under concurrent writes or a temporary wall-clock rollback. + */ + private static final AtomicLong MESSAGE_ID_GENERATOR = + new AtomicLong(System.currentTimeMillis() * 1_000_000L); + + protected static long nextMessageId() { + return MESSAGE_ID_GENERATOR.updateAndGet( + previous -> Math.max(previous + 1, System.currentTimeMillis() * 1_000_000L)); + } + /** * Timestamp represents the time of the current data point, in milliseconds. *
diff --git a/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/MeasureWrite.java b/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/MeasureWrite.java
index 477e21ef6501..e01582ca38e2 100644
--- a/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/MeasureWrite.java
+++ b/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/MeasureWrite.java
@@ -88,8 +88,10 @@ protected BanyandbMeasure.WriteRequest build(BanyandbCommon.Metadata metadata) {
}
builder.setDataPointSpec(datapointValueSpecBuilder);
+ long messageId = nextMessageId();
+ datapointValueBuilder.setVersion(messageId);
builder.setDataPoint(datapointValueBuilder);
- builder.setMessageId(System.nanoTime());
+ builder.setMessageId(messageId);
return builder.build();
}
@@ -122,8 +124,10 @@ protected BanyandbMeasure.WriteRequest buildValues() {
datapointValueBuilder.addFields(fieldEntry.getValue().serialize());
}
+ long messageId = nextMessageId();
+ datapointValueBuilder.setVersion(messageId);
builder.setDataPoint(datapointValueBuilder);
- builder.setMessageId(System.nanoTime());
+ builder.setMessageId(messageId);
return builder.build();
}
diff --git a/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/StreamWrite.java b/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/StreamWrite.java
index 91f80c0b1b51..3310065d875c 100644
--- a/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/StreamWrite.java
+++ b/oap-server/server-library/library-banyandb-client/src/main/java/org/apache/skywalking/library/banyandb/v1/client/StreamWrite.java
@@ -93,7 +93,7 @@ protected BanyandbStream.WriteRequest build(BanyandbCommon.Metadata metadata) {
}
builder.setElement(elemValBuilder);
- builder.setMessageId(System.nanoTime());
+ builder.setMessageId(nextMessageId());
return builder.build();
}
@@ -125,7 +125,7 @@ protected BanyandbStream.WriteRequest buildValues() {
}
builder.setElement(elemValBuilder);
- builder.setMessageId(System.nanoTime());
+ builder.setMessageId(nextMessageId());
return builder.build();
}
diff --git a/oap-server/server-library/library-banyandb-client/src/test/java/org/apache/skywalking/library/banyandb/v1/client/MeasureWriteMessageIdTest.java b/oap-server/server-library/library-banyandb-client/src/test/java/org/apache/skywalking/library/banyandb/v1/client/MeasureWriteMessageIdTest.java
new file mode 100644
index 000000000000..a2ca5994e21c
--- /dev/null
+++ b/oap-server/server-library/library-banyandb-client/src/test/java/org/apache/skywalking/library/banyandb/v1/client/MeasureWriteMessageIdTest.java
@@ -0,0 +1,101 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.skywalking.library.banyandb.v1.client;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import org.apache.skywalking.banyandb.common.v1.BanyandbCommon;
+import org.apache.skywalking.banyandb.measure.v1.BanyandbMeasure;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class MeasureWriteMessageIdTest {
+
+ private static final BanyandbCommon.Metadata METADATA =
+ BanyandbCommon.Metadata.newBuilder()
+ .setGroup("test-group")
+ .setName("test-measure")
+ .build();
+
+ @Test
+ public void messageIdMatchesDataPointVersion() {
+ final MeasureWrite write = new MeasureWrite(METADATA, System.currentTimeMillis());
+ final BanyandbMeasure.WriteRequest request = write.build();
+ assertEquals(request.getMessageId(), request.getDataPoint().getVersion());
+ assertTrue(request.getMessageId() > 0);
+ }
+
+ @Test
+ public void messageIdMatchesDataPointVersionInBuildValues() {
+ final MeasureWrite write = new MeasureWrite(METADATA, System.currentTimeMillis());
+ final BanyandbMeasure.WriteRequest request = write.buildOnlyValues();
+ assertEquals(request.getMessageId(), request.getDataPoint().getVersion());
+ assertTrue(request.getMessageId() > 0);
+ }
+
+ @Test
+ public void messageIdsAreStrictlyIncreasing() {
+ final List