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 ids = new ArrayList<>(); + for (int i = 0; i < 10_000; i++) { + ids.add(AbstractWrite.nextMessageId()); + } + for (int i = 1; i < ids.size(); i++) { + assertTrue(ids.get(i) > ids.get(i - 1), "message ids must be strictly increasing"); + } + } + + @Test + public void messageIdsAreUniqueUnderConcurrency() throws InterruptedException { + final int threads = 8; + final int perThread = 2_000; + final ExecutorService executor = Executors.newFixedThreadPool(threads); + final CountDownLatch start = new CountDownLatch(1); + final Set seen = ConcurrentHashMap.newKeySet(); + final CountDownLatch done = new CountDownLatch(threads); + try { + for (int t = 0; t < threads; t++) { + executor.submit(() -> { + try { + start.await(); + for (int i = 0; i < perThread; i++) { + seen.add(AbstractWrite.nextMessageId()); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + done.countDown(); + } + }); + } + start.countDown(); + assertTrue(done.await(30, TimeUnit.SECONDS), "concurrent writes timed out"); + } finally { + executor.shutdownNow(); + } + assertEquals(threads * perThread, seen.size(), "all message ids must be unique"); + } +} \ No newline at end of file