This is an automated email from the ASF dual-hosted git repository.
Caideyipi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new dc2e7c861b6 [Pipe] Optimize OPC UA TsFile last-value transfer (#18433)
dc2e7c861b6 is described below
commit dc2e7c861b62b049b1b6cc26b9eb2fd7d3245a0d
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 18 12:13:38 2026 +0800
[Pipe] Optimize OPC UA TsFile last-value transfer (#18433)
---
.../db/pipe/sink/protocol/opcua/OpcUaSink.java | 165 ++++++++
.../protocol/opcua/client/IoTDBOpcUaClient.java | 39 +-
.../sink/protocol/opcua/server/OpcUaNameSpace.java | 89 ++++-
.../opcua/OpcUaSinkTsFileMetadataTest.java | 167 +++++++++
.../opcua/OpcUaSinkTsFilePerformanceTest.java | 415 +++++++++++++++++++++
.../opcua/client/IoTDBOpcUaClientTest.java | 38 ++
.../opcua/server/OpcUaNameSpaceMetadataTest.java | 167 +++++++++
7 files changed, 1075 insertions(+), 5 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java
index 1263603bd49..af00b88ffa6 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java
@@ -25,6 +25,7 @@ import org.apache.iotdb.db.conf.IoTDBConfig;
import org.apache.iotdb.db.i18n.DataNodePipeMessages;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeInsertNodeTabletInsertionEvent;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
+import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import org.apache.iotdb.db.pipe.sink.protocol.opcua.client.ClientRunner;
import org.apache.iotdb.db.pipe.sink.protocol.opcua.client.IoTDBOpcUaClient;
import org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace;
@@ -39,11 +40,21 @@ import
org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
+import org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
import org.apache.iotdb.pipe.api.exception.PipeException;
import org.apache.tsfile.common.conf.TSFileConfig;
+import org.apache.tsfile.common.constant.TsFileConstant;
+import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.file.metadata.TimeseriesMetadata;
+import org.apache.tsfile.read.TimeValuePair;
+import org.apache.tsfile.read.TsFileSequenceReader;
+import org.apache.tsfile.read.reader.TsFileLastReader;
import org.apache.tsfile.utils.Pair;
import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.IMeasurementSchema;
+import org.apache.tsfile.write.schema.MeasurementSchema;
import org.eclipse.milo.opcua.sdk.client.identity.AnonymousProvider;
import org.eclipse.milo.opcua.sdk.client.identity.IdentityProvider;
import org.eclipse.milo.opcua.sdk.client.identity.UsernameProvider;
@@ -56,7 +67,11 @@ import org.slf4j.LoggerFactory;
import javax.annotation.Nullable;
import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
import java.util.Arrays;
+import java.util.LinkedHashMap;
+import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
@@ -497,6 +512,156 @@ public class OpcUaSink implements PipeConnector {
// Server side, do nothing
}
+ @Override
+ public void transfer(final TsFileInsertionEvent tsFileInsertionEvent) throws
Exception {
+ if (!shouldTransferTsFileByMetadata(tsFileInsertionEvent)) {
+ PipeConnector.super.transfer(tsFileInsertionEvent);
+ return;
+ }
+
+ final PipeTsFileInsertionEvent pipeTsFileInsertionEvent =
+ (PipeTsFileInsertionEvent) tsFileInsertionEvent;
+ boolean delegatedToTabletTransfer = false;
+ try {
+ if (transferTsFileByMetadata(pipeTsFileInsertionEvent)
+ == TsFileTransferResult.FALLBACK_TO_TABLETS) {
+ delegatedToTabletTransfer = true;
+ PipeConnector.super.transfer(tsFileInsertionEvent);
+ }
+ } finally {
+ // PipeConnector.transfer(TsFileInsertionEvent) closes the event itself
when it is used as a
+ // fallback. Keep the ownership here for the metadata fast path and
exceptional exits.
+ if (!delegatedToTabletTransfer) {
+ tsFileInsertionEvent.close();
+ }
+ }
+ }
+
+ private boolean shouldTransferTsFileByMetadata(final TsFileInsertionEvent
tsFileInsertionEvent) {
+ if (!isClientServerModel || !(tsFileInsertionEvent instanceof
PipeTsFileInsertionEvent)) {
+ return false;
+ }
+
+ final PipeTsFileInsertionEvent pipeTsFileInsertionEvent =
+ (PipeTsFileInsertionEvent) tsFileInsertionEvent;
+ // Metadata contains the unfiltered last value. Deletions, path/time
filters, and privilege
+ // filtering must use the normal parser so that the sink observes exactly
the event payload.
+ return !pipeTsFileInsertionEvent.isWithMod()
+ && !pipeTsFileInsertionEvent.shouldParseTimeOrPattern()
+ && !pipeTsFileInsertionEvent.shouldParse4Privilege();
+ }
+
+ private TsFileTransferResult transferTsFileByMetadata(
+ final PipeTsFileInsertionEvent pipeTsFileInsertionEvent) throws
Exception {
+ if
(!pipeTsFileInsertionEvent.increaseReferenceCount(OpcUaSink.class.getName())) {
+ return TsFileTransferResult.SKIPPED;
+ }
+
+ try {
+ if (!pipeTsFileInsertionEvent.waitForTsFileClose()) {
+ return TsFileTransferResult.SKIPPED;
+ }
+
+ final Map<IDeviceID, List<Pair<IMeasurementSchema, TimeValuePair>>>
deviceLastValues;
+ try {
+ deviceLastValues =
readLastValues(pipeTsFileInsertionEvent.getTsFile());
+ } catch (final Exception e) {
+ // Keep the parser as a compatibility fallback when the TsFile
metadata cannot be read.
+ return TsFileTransferResult.FALLBACK_TO_TABLETS;
+ }
+ if (Objects.isNull(deviceLastValues)) {
+ return TsFileTransferResult.FALLBACK_TO_TABLETS;
+ }
+
+ final boolean isTableModel =
pipeTsFileInsertionEvent.isTableModelEvent();
+ if (Objects.nonNull(nameSpace)) {
+ for (final Map.Entry<IDeviceID, List<Pair<IMeasurementSchema,
TimeValuePair>>> entry :
+ deviceLastValues.entrySet()) {
+ nameSpace.transferLastValues(entry.getKey(), entry.getValue(),
isTableModel, this);
+ }
+ } else if (Objects.nonNull(client)) {
+ // Batch all devices into the same OPC UA write so that many-device
TsFiles do not incur one
+ // network round trip per device.
+ client.transferLastValues(deviceLastValues, isTableModel, this);
+ } else {
+ throw new
PipeException(DataNodePipeMessages.NO_OPC_CLIENT_OR_SERVER_IS_SPECIFIED);
+ }
+ return TsFileTransferResult.TRANSFERRED;
+ } finally {
+
pipeTsFileInsertionEvent.decreaseReferenceCount(OpcUaSink.class.getName(),
false);
+ }
+ }
+
+ static @Nullable Map<IDeviceID, List<Pair<IMeasurementSchema,
TimeValuePair>>> readLastValues(
+ final File tsFile) throws Exception {
+ final Map<IDeviceID, Map<String, TSDataType>> deviceToTimeseriesDataTypes =
+ readTimeseriesDataTypes(tsFile);
+ final long expectedTimeseriesCount =
+
deviceToTimeseriesDataTypes.values().stream().mapToLong(Map::size).sum();
+ long actualTimeseriesCount = 0;
+ final Map<IDeviceID, List<Pair<IMeasurementSchema, TimeValuePair>>>
deviceLastValues =
+ new LinkedHashMap<>();
+ // Disable asynchronous IO here. The sink already runs in a pipe worker
and a synchronous
+ // reader avoids leaving a background task behind when the event is
cancelled or falls back to
+ // tablet parsing.
+ try (final TsFileLastReader lastReader = new
TsFileLastReader(tsFile.getPath(), false, false)) {
+ while (lastReader.hasNext()) {
+ final Pair<IDeviceID, List<Pair<String, TimeValuePair>>>
deviceLastValue =
+ lastReader.next();
+ final Map<String, TSDataType> timeseriesDataTypes =
+ deviceToTimeseriesDataTypes.get(deviceLastValue.getLeft());
+ if (Objects.isNull(timeseriesDataTypes)) {
+ return null;
+ }
+
+ final List<Pair<IMeasurementSchema, TimeValuePair>> typedLastValues =
+ deviceLastValues.computeIfAbsent(deviceLastValue.getLeft(), key ->
new ArrayList<>());
+ for (final Pair<String, TimeValuePair> lastValue :
deviceLastValue.getRight()) {
+ ++actualTimeseriesCount;
+ final TSDataType dataType =
timeseriesDataTypes.get(lastValue.getLeft());
+ if (Objects.isNull(dataType)) {
+ return null;
+ }
+ if (!TsFileConstant.TIME_COLUMN_ID.equals(lastValue.getLeft())) {
+ typedLastValues.add(
+ new Pair<>(
+ new MeasurementSchema(lastValue.getLeft(), dataType),
lastValue.getRight()));
+ }
+ }
+ }
+ }
+
+ // TsFileLastReader logs and suppresses IOExceptions from
Iterator#hasNext. Comparing against an
+ // independently read metadata count prevents a truncated result from
being treated as EOF.
+ if (actualTimeseriesCount != expectedTimeseriesCount) {
+ return null;
+ }
+ return deviceLastValues;
+ }
+
+ private static Map<IDeviceID, Map<String, TSDataType>>
readTimeseriesDataTypes(final File tsFile)
+ throws IOException {
+ try (final TsFileSequenceReader sequenceReader = new
TsFileSequenceReader(tsFile.getPath())) {
+ final Map<IDeviceID, Map<String, TSDataType>>
deviceToTimeseriesDataTypes =
+ new LinkedHashMap<>();
+ for (final Map.Entry<IDeviceID, List<TimeseriesMetadata>> entry :
+ sequenceReader.getAllTimeseriesMetadata(false).entrySet()) {
+ final Map<String, TSDataType> timeseriesDataTypes = new
LinkedHashMap<>();
+ for (final TimeseriesMetadata metadata : entry.getValue()) {
+ timeseriesDataTypes.put(metadata.getMeasurementId(),
metadata.getTsDataType());
+ }
+ deviceToTimeseriesDataTypes.put(entry.getKey(), timeseriesDataTypes);
+ }
+ return deviceToTimeseriesDataTypes;
+ }
+ }
+
+ private enum TsFileTransferResult {
+ TRANSFERRED,
+ SKIPPED,
+ FALLBACK_TO_TABLETS
+ }
+
@Override
public void transfer(final Event event) throws Exception {
// Do nothing when receive heartbeat or other events
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClient.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClient.java
index 8c6c30dfeba..83fc94e2ea2 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClient.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClient.java
@@ -28,6 +28,9 @@ import org.apache.iotdb.pipe.api.exception.PipeException;
import org.apache.tsfile.common.constant.TsFileConstant;
import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.read.TimeValuePair;
+import org.apache.tsfile.utils.Pair;
import org.apache.tsfile.write.record.Tablet;
import org.apache.tsfile.write.schema.IMeasurementSchema;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
@@ -63,6 +66,7 @@ import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
+import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.ExecutionException;
@@ -123,6 +127,26 @@ public class IoTDBOpcUaClient {
tablet, false, sink, this::transferTabletRowForClientServerModel);
}
+ public void transferLastValues(
+ final Map<IDeviceID, List<Pair<IMeasurementSchema, TimeValuePair>>>
deviceLastValues,
+ final boolean isTableModel,
+ final OpcUaSink sink)
+ throws Exception {
+ final List<OpcUaWriteRequest> writeRequests = new ArrayList<>();
+ for (final Map.Entry<IDeviceID, List<Pair<IMeasurementSchema,
TimeValuePair>>> entry :
+ deviceLastValues.entrySet()) {
+ OpcUaNameSpace.transferLastValues(
+ entry.getKey(),
+ entry.getValue(),
+ isTableModel,
+ sink,
+ (segments, measurementSchemas, timestamps, values, currentSink) ->
+ collectWriteRequests(
+ segments, measurementSchemas, timestamps, values,
currentSink, writeRequests));
+ }
+ writeValues(writeRequests);
+ }
+
private void transferTabletRowForClientServerModel(
final String[] segments,
final List<IMeasurementSchema> measurementSchemas,
@@ -130,11 +154,22 @@ public class IoTDBOpcUaClient {
final List<Object> values,
final OpcUaSink sink)
throws Exception {
+ final List<OpcUaWriteRequest> writeRequests = new ArrayList<>();
+ collectWriteRequests(segments, measurementSchemas, timestamps, values,
sink, writeRequests);
+ writeValues(writeRequests);
+ }
+
+ private void collectWriteRequests(
+ final String[] segments,
+ final List<IMeasurementSchema> measurementSchemas,
+ final List<Long> timestamps,
+ final List<Object> values,
+ final OpcUaSink sink,
+ final List<OpcUaWriteRequest> writeRequests) {
StatusCode currentQuality = sink.getDefaultQuality();
Object value = null;
long timestamp = 0;
NodeId opcDataType = null;
- final List<OpcUaWriteRequest> writeRequests = new ArrayList<>();
for (int i = 0; i < measurementSchemas.size(); ++i) {
if (Objects.isNull(values.get(i))) {
@@ -177,8 +212,6 @@ public class IoTDBOpcUaClient {
writeRequests.add(
new OpcUaWriteRequest(value, timestamp, opcDataType, currentQuality,
segments, null));
}
-
- writeValues(writeRequests);
}
private void writeValues(final List<OpcUaWriteRequest> writeRequests) throws
Exception {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java
index 1d6262c2c6a..421ba27f77f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java
@@ -33,7 +33,11 @@ import org.apache.iotdb.pipe.api.event.Event;
import org.apache.tsfile.common.constant.TsFileConstant;
import org.apache.tsfile.enums.ColumnCategory;
import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.read.TimeValuePair;
import org.apache.tsfile.utils.Binary;
+import org.apache.tsfile.utils.DateUtils;
+import org.apache.tsfile.utils.Pair;
import org.apache.tsfile.write.UnSupportedDataTypeException;
import org.apache.tsfile.write.record.Tablet;
import org.apache.tsfile.write.schema.IMeasurementSchema;
@@ -63,12 +67,12 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.file.Paths;
-import java.sql.Date;
import java.time.LocalDate;
import java.time.ZoneId;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.Date;
import java.util.List;
import java.util.Objects;
import java.util.Set;
@@ -130,6 +134,85 @@ public class OpcUaNameSpace extends
ManagedNamespaceWithLifecycle {
}
}
+ /**
+ * Transfers the last value of every measurement in a TsFile device without
materializing a {@link
+ * Tablet}. The TsFile last-value reader obtains the values from metadata
(and reads only the last
+ * chunk when a data type does not keep a value in statistics).
+ */
+ public void transferLastValues(
+ final IDeviceID deviceID,
+ final List<Pair<IMeasurementSchema, TimeValuePair>> lastValues,
+ final boolean isTableModel,
+ final OpcUaSink sink)
+ throws Exception {
+ transferLastValues(
+ deviceID, lastValues, isTableModel, sink,
this::transferTabletRowForClientServerModel);
+ }
+
+ public static void transferLastValues(
+ final IDeviceID deviceID,
+ final List<Pair<IMeasurementSchema, TimeValuePair>> lastValues,
+ final boolean isTableModel,
+ final OpcUaSink sink,
+ final TabletRowConsumer consumer)
+ throws Exception {
+ final String[] segments;
+ if (!isTableModel) {
+ // IDeviceID may compact multiple tree nodes into one segment. Keep the
same node layout as
+ // the Tablet path, which splits the complete device path.
+ segments = deviceID.toString().split("\\.");
+ } else {
+ final Object[] deviceSegments = deviceID.getSegments();
+ segments = new String[deviceSegments.length + 1];
+ segments[0] = sink.getDatabaseName();
+ for (int i = 0; i < deviceSegments.length; ++i) {
+ segments[i + 1] =
+ Objects.isNull(deviceSegments[i])
+ ? sink.getPlaceHolder4NullTag()
+ : String.valueOf(deviceSegments[i]);
+ }
+ }
+
+ final List<IMeasurementSchema> schemas = new
ArrayList<>(lastValues.size());
+ final List<Long> timestamps = new ArrayList<>(lastValues.size());
+ final List<Object> values = new ArrayList<>(lastValues.size());
+ for (final Pair<IMeasurementSchema, TimeValuePair> lastValue : lastValues)
{
+ if (Objects.isNull(lastValue)
+ || Objects.isNull(lastValue.getLeft())
+ || Objects.isNull(lastValue.getLeft().getMeasurementName())
+ ||
TsFileConstant.TIME_COLUMN_ID.equals(lastValue.getLeft().getMeasurementName())
+ || Objects.isNull(lastValue.getRight())
+ || Objects.isNull(lastValue.getRight().getValue())) {
+ continue;
+ }
+
+ final TimeValuePair timeValuePair = lastValue.getRight();
+ final TSDataType dataType = lastValue.getLeft().getType();
+ schemas.add(lastValue.getLeft());
+ timestamps.add(timeValuePair.getTimestamp());
+ values.add(getObjectValue4Opc(timeValuePair, dataType));
+ }
+
+ if (!schemas.isEmpty()) {
+ consumer.accept(segments, schemas, timestamps, values, sink);
+ }
+ }
+
+ private static Object getObjectValue4Opc(
+ final TimeValuePair timeValuePair, final TSDataType dataType) {
+ final Object value = timeValuePair.getValue().getValue();
+ return switch (dataType) {
+ case DATE ->
+ new DateTime(new Date(DateUtils.parseIntToDate(((Number)
value).intValue()).getTime()));
+ case TIMESTAMP -> new DateTime(timestampToUtc(((Number)
value).longValue()));
+ case TEXT, BLOB, STRING -> value instanceof Binary ? value.toString() :
String.valueOf(value);
+ case BOOLEAN, INT32, INT64, FLOAT, DOUBLE -> value;
+ case VECTOR, OBJECT, UNKNOWN ->
+ throw new UnSupportedDataTypeException(
+ DataNodePipeMessages.UNSUPPORTED_DATATYPE + dataType);
+ };
+ }
+
public static void transferTabletForClientServerModel(
final Tablet tablet,
final boolean isTableModel,
@@ -393,7 +476,9 @@ public class OpcUaNameSpace extends
ManagedNamespaceWithLifecycle {
case INT32:
return ((int[]) column)[rowIndex];
case DATE:
- return new DateTime(Date.valueOf(((LocalDate[]) column)[rowIndex]));
+ return new DateTime(
+ Date.from(
+ ((LocalDate[])
column)[rowIndex].atStartOfDay(ZoneId.systemDefault()).toInstant()));
case INT64:
return ((long[]) column)[rowIndex];
case TIMESTAMP:
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSinkTsFileMetadataTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSinkTsFileMetadataTest.java
new file mode 100644
index 00000000000..ad26debe416
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSinkTsFileMetadataTest.java
@@ -0,0 +1,167 @@
+/*
+ * 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.iotdb.db.pipe.sink.protocol.opcua;
+
+import org.apache.tsfile.common.conf.TSFileConfig;
+import org.apache.tsfile.common.constant.TsFileConstant;
+import org.apache.tsfile.enums.ColumnCategory;
+import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.file.metadata.TableSchema;
+import org.apache.tsfile.read.TimeValuePair;
+import org.apache.tsfile.utils.Binary;
+import org.apache.tsfile.utils.Pair;
+import org.apache.tsfile.write.TsFileWriter;
+import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.IMeasurementSchema;
+import org.apache.tsfile.write.schema.MeasurementSchema;
+import org.junit.Assert;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+
+import java.io.File;
+import java.util.Arrays;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+public class OpcUaSinkTsFileMetadataTest {
+
+ @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();
+
+ @Test
+ public void testReadLastValuesFromTreeTsFile() throws Exception {
+ final File tsFile = new File(temporaryFolder.getRoot(), "tree.tsfile");
+ final String device = "root.sg.d1";
+ final List<IMeasurementSchema> schemas =
+ Arrays.asList(
+ new MeasurementSchema("s1", TSDataType.INT64),
+ new MeasurementSchema("blob", TSDataType.BLOB));
+ final Tablet tablet = new Tablet(device, schemas, 3);
+ for (int i = 0; i < 3; ++i) {
+ tablet.addTimestamp(i, i + 1L);
+ }
+ tablet.addValue("s1", 0, 10L);
+ tablet.addValue("s1", 1, 20L);
+ tablet.addValue("s1", 2, null);
+ tablet.addValue("blob", 0, new Binary("old", TSFileConfig.STRING_CHARSET));
+ tablet.addValue("blob", 1, null);
+ tablet.addValue("blob", 2, new Binary("last",
TSFileConfig.STRING_CHARSET));
+
+ try (final TsFileWriter writer = new TsFileWriter(tsFile)) {
+ for (final IMeasurementSchema schema : schemas) {
+ writer.registerTimeseries(device, schema);
+ }
+ writer.writeTree(tablet);
+ }
+
+ final Map<IDeviceID, List<Pair<IMeasurementSchema, TimeValuePair>>>
deviceLastValues =
+ OpcUaSink.readLastValues(tsFile);
+ Assert.assertEquals(1, deviceLastValues.size());
+ final Map<String, TimeValuePair> lastValues =
+ toMeasurementMap(deviceLastValues.values().iterator().next());
+
+ Assert.assertFalse(lastValues.containsKey(TsFileConstant.TIME_COLUMN_ID));
+ assertLongLastValue(lastValues.get("s1"), 2L, 20L);
+ assertBinaryLastValue(lastValues.get("blob"), 3L, "last");
+ }
+
+ @Test
+ public void testReadLastValuesFromTableTsFile() throws Exception {
+ final File tsFile = new File(temporaryFolder.getRoot(), "table.tsfile");
+ final List<String> columnNames = Arrays.asList("tag", "s1", "blob",
"timestamp");
+ final List<TSDataType> dataTypes =
+ Arrays.asList(TSDataType.STRING, TSDataType.INT64, TSDataType.BLOB,
TSDataType.TIMESTAMP);
+ final List<ColumnCategory> columnCategories =
+ Arrays.asList(
+ ColumnCategory.TAG, ColumnCategory.FIELD, ColumnCategory.FIELD,
ColumnCategory.FIELD);
+ final List<IMeasurementSchema> schemas =
+ Arrays.asList(
+ new MeasurementSchema("tag", TSDataType.STRING),
+ new MeasurementSchema("s1", TSDataType.INT64),
+ new MeasurementSchema("blob", TSDataType.BLOB),
+ new MeasurementSchema("timestamp", TSDataType.TIMESTAMP));
+ final Tablet tablet = new Tablet("table", columnNames, dataTypes,
columnCategories, 3);
+ for (int i = 0; i < 3; ++i) {
+ tablet.addTimestamp(i, i + 1L);
+ tablet.addValue(i, 0, "tag-value");
+ }
+ tablet.addValue(0, 1, 10L);
+ tablet.addValue(1, 1, 20L);
+ tablet.addValue("s1", 2, null);
+ tablet.addValue("blob", 0, new Binary("old", TSFileConfig.STRING_CHARSET));
+ tablet.addValue("blob", 1, null);
+ tablet.addValue("blob", 2, new Binary("last",
TSFileConfig.STRING_CHARSET));
+ tablet.addValue(0, 3, 1_700_000_000_000L);
+ tablet.addValue(1, 3, 1_700_000_001_000L);
+ tablet.addValue(2, 3, 1_700_000_002_000L);
+
+ try (final TsFileWriter writer = new TsFileWriter(tsFile)) {
+ writer.registerTableSchema(new TableSchema("table", schemas,
columnCategories));
+ writer.writeTable(tablet);
+ }
+
+ final Map<IDeviceID, List<Pair<IMeasurementSchema, TimeValuePair>>>
deviceLastValues =
+ OpcUaSink.readLastValues(tsFile);
+ Assert.assertEquals(1, deviceLastValues.size());
+ final Map.Entry<IDeviceID, List<Pair<IMeasurementSchema, TimeValuePair>>>
entry =
+ deviceLastValues.entrySet().iterator().next();
+ Assert.assertArrayEquals(new Object[] {"table", "tag-value"},
entry.getKey().getSegments());
+
+ final Map<String, TimeValuePair> lastValues =
toMeasurementMap(entry.getValue());
+ Assert.assertFalse(lastValues.containsKey("tag"));
+ Assert.assertFalse(lastValues.containsKey(TsFileConstant.TIME_COLUMN_ID));
+ assertLongLastValue(lastValues.get("s1"), 2L, 20L);
+ assertBinaryLastValue(lastValues.get("blob"), 3L, "last");
+ assertLongLastValue(lastValues.get("timestamp"), 3L, 1_700_000_002_000L);
+ Assert.assertEquals(TSDataType.TIMESTAMP, getSchema(entry.getValue(),
"timestamp").getType());
+ }
+
+ private static Map<String, TimeValuePair> toMeasurementMap(
+ final List<Pair<IMeasurementSchema, TimeValuePair>> lastValues) {
+ final Map<String, TimeValuePair> result = new LinkedHashMap<>();
+ lastValues.forEach(
+ lastValue -> result.put(lastValue.getLeft().getMeasurementName(),
lastValue.getRight()));
+ return result;
+ }
+
+ private static IMeasurementSchema getSchema(
+ final List<Pair<IMeasurementSchema, TimeValuePair>> lastValues, final
String measurement) {
+ return lastValues.stream()
+ .map(Pair::getLeft)
+ .filter(schema -> measurement.equals(schema.getMeasurementName()))
+ .findFirst()
+ .orElseThrow(AssertionError::new);
+ }
+
+ private static void assertLongLastValue(
+ final TimeValuePair lastValue, final long timestamp, final long value) {
+ Assert.assertNotNull(lastValue);
+ Assert.assertEquals(timestamp, lastValue.getTimestamp());
+ Assert.assertEquals(value, lastValue.getValue().getLong());
+ }
+
+ private static void assertBinaryLastValue(
+ final TimeValuePair lastValue, final long timestamp, final String value)
{
+ Assert.assertNotNull(lastValue);
+ Assert.assertEquals(timestamp, lastValue.getTimestamp());
+ Assert.assertEquals(value, lastValue.getValue().getBinary().toString());
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSinkTsFilePerformanceTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSinkTsFilePerformanceTest.java
new file mode 100644
index 00000000000..59fdd76f73a
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSinkTsFilePerformanceTest.java
@@ -0,0 +1,415 @@
+/*
+ * 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.iotdb.db.pipe.sink.protocol.opcua;
+
+import org.apache.iotdb.commons.conf.CommonDescriptor;
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.pipe.config.PipeConfig;
+import org.apache.iotdb.commons.pipe.datastructure.pattern.PrefixTreePattern;
+import
org.apache.iotdb.db.pipe.event.common.tsfile.parser.scan.TsFileInsertionEventScanParser;
+import org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace;
+import
org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace.TabletRowConsumer;
+import org.apache.iotdb.db.utils.ManualPerformanceTestUtils;
+import org.apache.iotdb.db.utils.ManualPerformanceTestUtils.Measurement;
+import org.apache.iotdb.db.utils.ManualPerformanceTestUtils.Summary;
+
+import org.apache.tsfile.common.conf.TSFileConfig;
+import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.file.metadata.enums.CompressionType;
+import org.apache.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.tsfile.read.TimeValuePair;
+import org.apache.tsfile.utils.Binary;
+import org.apache.tsfile.utils.Pair;
+import org.apache.tsfile.write.TsFileWriter;
+import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.IMeasurementSchema;
+import org.apache.tsfile.write.schema.MeasurementSchema;
+import org.junit.Assert;
+import org.junit.Assume;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+
+import java.io.File;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.Objects;
+
+/**
+ * Manual benchmark for the OPC UA client-server TsFile paths. Enable it with
{@code
+ * -Diotdb.opcua.tsfile.last-value.perf.enabled=true}; the remaining
properties below tune the
+ * generated file and measurement rounds.
+ */
+public class OpcUaSinkTsFilePerformanceTest {
+
+ private static final String ENABLED_PROPERTY =
"iotdb.opcua.tsfile.last-value.perf.enabled";
+ private static final String DEVICE_COUNT_PROPERTY =
+ "iotdb.opcua.tsfile.last-value.perf.device.count";
+ private static final String MEASUREMENT_COUNT_PROPERTY =
+ "iotdb.opcua.tsfile.last-value.perf.measurement.count";
+ private static final String BLOB_MEASUREMENT_COUNT_PROPERTY =
+ "iotdb.opcua.tsfile.last-value.perf.blob.measurement.count";
+ private static final String ROW_COUNT_PROPERTY =
"iotdb.opcua.tsfile.last-value.perf.row.count";
+ private static final String TABLET_ROW_COUNT_PROPERTY =
+ "iotdb.opcua.tsfile.last-value.perf.tablet.row.count";
+ private static final String WARMUP_ITERATIONS_PROPERTY =
+ "iotdb.opcua.tsfile.last-value.perf.warmup.iterations";
+ private static final String ITERATIONS_PROPERTY =
"iotdb.opcua.tsfile.last-value.perf.iterations";
+ private static final String ROUNDS_PROPERTY =
"iotdb.opcua.tsfile.last-value.perf.rounds";
+
+ private static volatile long benchmarkBlackhole;
+
+ @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();
+
+ @Test
+ public void compareTabletAndMetadataLastValuePaths() throws Exception {
+ Assume.assumeTrue(
+ String.format(
+ Locale.ROOT,
+ "Manual performance UT. Enable with -D%s=true; use the
benchmark-specific properties to tune its workload.",
+ ENABLED_PROPERTY),
+ Boolean.getBoolean(ENABLED_PROPERTY));
+ Assume.assumeTrue(
+ "Current-thread CPU time and allocation metrics are required.",
+ ManualPerformanceTestUtils.enableThreadMetrics());
+
+ final int deviceCount = Integer.getInteger(DEVICE_COUNT_PROPERTY, 2);
+ final int measurementCount =
Integer.getInteger(MEASUREMENT_COUNT_PROPERTY, 32);
+ final int blobMeasurementCount =
Integer.getInteger(BLOB_MEASUREMENT_COUNT_PROPERTY, 1);
+ final int rowCount = Integer.getInteger(ROW_COUNT_PROPERTY, 50_000);
+ final int tabletRowCount = Integer.getInteger(TABLET_ROW_COUNT_PROPERTY,
1024);
+ final int warmupIterations =
Integer.getInteger(WARMUP_ITERATIONS_PROPERTY, 1);
+ final int iterations = Integer.getInteger(ITERATIONS_PROPERTY, 1);
+ final int rounds = Integer.getInteger(ROUNDS_PROPERTY, 5);
+ assertPositive(
+ deviceCount,
+ measurementCount,
+ rowCount,
+ tabletRowCount,
+ warmupIterations,
+ iterations,
+ rounds);
+ Assert.assertTrue(blobMeasurementCount >= 0);
+ Assert.assertTrue(blobMeasurementCount <= measurementCount);
+
+ final File tsFile = new File(temporaryFolder.getRoot(),
"opcua-last-value-performance.tsfile");
+ generateAlignedTreeTsFile(
+ tsFile, deviceCount, measurementCount, blobMeasurementCount, rowCount,
tabletRowCount);
+
+ final boolean originalPipeMemoryManagementEnabled =
+ PipeConfig.getInstance().getPipeMemoryManagementEnabled();
+
CommonDescriptor.getInstance().getConfig().setPipeMemoryManagementEnabled(false);
+ try {
+ final Map<String, CapturedLastValue> tabletLastValues =
+ captureLastValues(consumer -> transferByTabletPath(tsFile,
consumer));
+ final Map<String, CapturedLastValue> metadataLastValues =
+ captureLastValues(consumer -> transferByMetadataPath(tsFile,
consumer));
+ Assert.assertEquals((long) deviceCount * measurementCount,
tabletLastValues.size());
+ Assert.assertEquals(tabletLastValues, metadataLastValues);
+
+ compare(
+ tsFile,
+ deviceCount,
+ measurementCount,
+ blobMeasurementCount,
+ rowCount,
+ warmupIterations,
+ iterations,
+ rounds);
+ } finally {
+ CommonDescriptor.getInstance()
+ .getConfig()
+ .setPipeMemoryManagementEnabled(originalPipeMemoryManagementEnabled);
+ }
+ }
+
+ private static void compare(
+ final File tsFile,
+ final int deviceCount,
+ final int measurementCount,
+ final int blobMeasurementCount,
+ final int rowCount,
+ final int warmupIterations,
+ final int iterations,
+ final int rounds) {
+ final Runnable tabletPath = () -> benchmark(() ->
transferByTabletPath(tsFile, null));
+ final Runnable metadataPath = () -> benchmark(() ->
transferByMetadataPath(tsFile, null));
+
+ for (int i = 0; i < warmupIterations; ++i) {
+ if ((i & 1) == 0) {
+ tabletPath.run();
+ metadataPath.run();
+ } else {
+ metadataPath.run();
+ tabletPath.run();
+ }
+ }
+
+ final Measurement[] tabletMeasurements = new Measurement[rounds];
+ final Measurement[] metadataMeasurements = new Measurement[rounds];
+ for (int round = 0; round < rounds; ++round) {
+ if ((round & 1) == 0) {
+ tabletMeasurements[round] =
ManualPerformanceTestUtils.measure(iterations, tabletPath);
+ metadataMeasurements[round] =
ManualPerformanceTestUtils.measure(iterations, metadataPath);
+ } else {
+ metadataMeasurements[round] =
ManualPerformanceTestUtils.measure(iterations, metadataPath);
+ tabletMeasurements[round] =
ManualPerformanceTestUtils.measure(iterations, tabletPath);
+ }
+ }
+
+ final Summary tabletSummary =
+ ManualPerformanceTestUtils.summarize(tabletMeasurements, iterations);
+ final Summary metadataSummary =
+ ManualPerformanceTestUtils.summarize(metadataMeasurements, iterations);
+ final long pointCount = (long) deviceCount * measurementCount * rowCount;
+ System.out.printf(
+ Locale.ROOT,
+ "%nOPC UA TsFile last-value benchmark: file=%.2f MiB, devices=%d,
measurements/device=%d, BLOB measurements/device=%d, rows/device=%d, points=%d,
warmups=%d, iterations/round=%d, rounds=%d%n",
+ tsFile.length() / 1024.0 / 1024.0,
+ deviceCount,
+ measurementCount,
+ blobMeasurementCount,
+ rowCount,
+ pointCount,
+ warmupIterations,
+ iterations,
+ rounds);
+ printSummary("tablet path", tabletSummary);
+ printSummary("metadata path", metadataSummary);
+ System.out.printf(
+ Locale.ROOT,
+ " change: CPU speedup=%.2fx, allocation reduction=%.1f%%, peak-heap
reduction=%.1f%%%n",
+ ratio(tabletSummary.getCpuNanosPerOperation(),
metadataSummary.getCpuNanosPerOperation()),
+ reduction(
+ tabletSummary.getAllocatedBytesPerOperation(),
+ metadataSummary.getAllocatedBytesPerOperation()),
+ reduction(tabletSummary.getPeakHeapDeltaBytes(),
metadataSummary.getPeakHeapDeltaBytes()));
+ }
+
+ private static void transferByTabletPath(
+ final File tsFile, final TabletRowConsumer suppliedConsumer) throws
Exception {
+ final BenchmarkConsumer benchmarkConsumer =
+ Objects.isNull(suppliedConsumer) ? new BenchmarkConsumer() : null;
+ final TabletRowConsumer consumer =
+ Objects.isNull(suppliedConsumer) ? benchmarkConsumer :
suppliedConsumer;
+ try (final TsFileInsertionEventScanParser parser =
+ new TsFileInsertionEventScanParser(
+ tsFile,
+ new PrefixTreePattern("root"),
+ Long.MIN_VALUE,
+ Long.MAX_VALUE,
+ null,
+ null,
+ false)) {
+ for (final Pair<Tablet, Boolean> tabletWithIsAligned :
parser.toTabletWithIsAligneds()) {
+ OpcUaNameSpace.transferTabletForClientServerModel(
+ tabletWithIsAligned.getLeft(), false, null, consumer);
+ }
+ }
+ if (Objects.nonNull(benchmarkConsumer)) {
+ benchmarkBlackhole = benchmarkConsumer.result();
+ }
+ }
+
+ private static void transferByMetadataPath(
+ final File tsFile, final TabletRowConsumer suppliedConsumer) throws
Exception {
+ final BenchmarkConsumer benchmarkConsumer =
+ Objects.isNull(suppliedConsumer) ? new BenchmarkConsumer() : null;
+ final TabletRowConsumer consumer =
+ Objects.isNull(suppliedConsumer) ? benchmarkConsumer :
suppliedConsumer;
+ for (final Map.Entry<IDeviceID, List<Pair<IMeasurementSchema,
TimeValuePair>>> entry :
+ OpcUaSink.readLastValues(tsFile).entrySet()) {
+ OpcUaNameSpace.transferLastValues(entry.getKey(), entry.getValue(),
false, null, consumer);
+ }
+ if (Objects.nonNull(benchmarkConsumer)) {
+ benchmarkBlackhole = benchmarkConsumer.result();
+ }
+ }
+
+ private static Map<String, CapturedLastValue> captureLastValues(
+ final ThrowingConsumerRunner runner) throws Exception {
+ final Map<String, CapturedLastValue> lastValues = new LinkedHashMap<>();
+ runner.run(
+ (segments, schemas, timestamps, values, sink) -> {
+ final String device = String.join(".", segments);
+ for (int i = 0; i < schemas.size(); ++i) {
+ lastValues.put(
+ device + "." + schemas.get(i).getMeasurementName(),
+ new CapturedLastValue(timestamps.get(i), values.get(i)));
+ }
+ });
+ return lastValues;
+ }
+
+ private static void benchmark(final ThrowingRunnable operation) {
+ try {
+ operation.run();
+ } catch (final Exception e) {
+ throw new AssertionError(e);
+ }
+ }
+
+ private static void generateAlignedTreeTsFile(
+ final File tsFile,
+ final int deviceCount,
+ final int measurementCount,
+ final int blobMeasurementCount,
+ final int rowCount,
+ final int tabletRowCount)
+ throws Exception {
+ final List<IMeasurementSchema> schemas = new ArrayList<>(measurementCount);
+ for (int measurement = 0; measurement < measurementCount; ++measurement) {
+ final TSDataType dataType =
+ measurement < blobMeasurementCount ? TSDataType.BLOB :
TSDataType.INT64;
+ schemas.add(
+ new MeasurementSchema(
+ "s" + measurement, dataType, TSEncoding.PLAIN,
CompressionType.LZ4));
+ }
+
+ try (final TsFileWriter writer = new TsFileWriter(tsFile)) {
+ for (int device = 0; device < deviceCount; ++device) {
+ final String deviceId = "root.opcua_perf.d" + device;
+ writer.registerAlignedTimeseries(new PartialPath(deviceId), schemas);
+ final Tablet tablet = new Tablet(deviceId, schemas, tabletRowCount);
+ for (int row = 0; row < rowCount; ++row) {
+ if (tablet.getRowSize() == tablet.getMaxRowNumber()) {
+ writer.writeAligned(tablet);
+ tablet.reset();
+ }
+
+ final int rowIndex = tablet.getRowSize();
+ tablet.addTimestamp(rowIndex, row);
+ for (int measurement = 0; measurement < measurementCount;
++measurement) {
+ if (measurement < blobMeasurementCount) {
+ tablet.addValue(
+ schemas.get(measurement).getMeasurementName(),
+ rowIndex,
+ new Binary(
+ "d" + device + "s" + measurement + "r" + row,
TSFileConfig.STRING_CHARSET));
+ } else {
+ tablet.addValue(
+ rowIndex,
+ measurement,
+ ((long) device * measurementCount + measurement) * rowCount
+ row);
+ }
+ }
+ }
+ if (tablet.getRowSize() > 0) {
+ writer.writeAligned(tablet);
+ }
+ }
+ }
+ }
+
+ private static void assertPositive(final int... values) {
+ for (final int value : values) {
+ Assert.assertTrue(value > 0);
+ }
+ }
+
+ private static void printSummary(final String label, final Summary summary) {
+ System.out.printf(
+ Locale.ROOT,
+ " %-13s CPU=%.3f ms/file, allocated=%.3f MiB/file, peak heap
delta=%.3f MiB%n",
+ label,
+ summary.getCpuNanosPerOperation() / 1_000_000.0,
+ summary.getAllocatedBytesPerOperation() / 1024.0 / 1024.0,
+ summary.getPeakHeapDeltaBytes() / 1024.0 / 1024.0);
+ }
+
+ private static double ratio(final double baseline, final double optimized) {
+ return optimized == 0 ? Double.POSITIVE_INFINITY : baseline / optimized;
+ }
+
+ private static double reduction(final double baseline, final double
optimized) {
+ return baseline == 0 ? 0 : (baseline - optimized) * 100.0 / baseline;
+ }
+
+ @FunctionalInterface
+ private interface ThrowingConsumerRunner {
+ void run(TabletRowConsumer consumer) throws Exception;
+ }
+
+ @FunctionalInterface
+ private interface ThrowingRunnable {
+ void run() throws Exception;
+ }
+
+ private static final class BenchmarkConsumer implements TabletRowConsumer {
+
+ private long hash = 1;
+ private long callbackCount;
+ private long valueCount;
+
+ @Override
+ public void accept(
+ final String[] segments,
+ final List<IMeasurementSchema> measurementSchemas,
+ final List<Long> timestamps,
+ final List<Object> values,
+ final OpcUaSink sink) {
+ ++callbackCount;
+ for (final String segment : segments) {
+ hash = 31 * hash + Objects.hashCode(segment);
+ }
+ for (int i = 0; i < measurementSchemas.size(); ++i) {
+ hash = 31 * hash +
measurementSchemas.get(i).getMeasurementName().hashCode();
+ hash = 31 * hash + Long.hashCode(timestamps.get(i));
+ hash = 31 * hash + Objects.hashCode(values.get(i));
+ ++valueCount;
+ }
+ }
+
+ private long result() {
+ return hash ^ callbackCount ^ valueCount;
+ }
+ }
+
+ private static final class CapturedLastValue {
+
+ private final long timestamp;
+ private final Object value;
+
+ private CapturedLastValue(final long timestamp, final Object value) {
+ this.timestamp = timestamp;
+ this.value = value;
+ }
+
+ @Override
+ public boolean equals(final Object obj) {
+ if (this == obj) {
+ return true;
+ }
+ if (!(obj instanceof CapturedLastValue)) {
+ return false;
+ }
+ final CapturedLastValue that = (CapturedLastValue) obj;
+ return timestamp == that.timestamp && Objects.equals(value, that.value);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(timestamp, value);
+ }
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java
index 5cb881939d4..8f8333143e3 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java
@@ -23,7 +23,12 @@ import
org.apache.iotdb.db.pipe.sink.protocol.opcua.OpcUaSink;
import org.apache.iotdb.pipe.api.exception.PipeException;
import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.read.TimeValuePair;
+import org.apache.tsfile.utils.Pair;
+import org.apache.tsfile.utils.TsPrimitiveType;
import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.IMeasurementSchema;
import org.apache.tsfile.write.schema.MeasurementSchema;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
import org.eclipse.milo.opcua.sdk.client.identity.AnonymousProvider;
@@ -42,7 +47,9 @@ import org.mockito.Mockito;
import java.util.Arrays;
import java.util.Collections;
+import java.util.LinkedHashMap;
import java.util.List;
+import java.util.Map;
import java.util.concurrent.CompletableFuture;
public class IoTDBOpcUaClientTest {
@@ -63,6 +70,30 @@ public class IoTDBOpcUaClientTest {
Mockito.argThat(listWithSize(2)));
}
+ @Test
+ public void testTransferLastValuesBatchesDevicesInOneRequest() throws
Exception {
+ final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
+ Mockito.when(miloClient.writeValuesAsync(Mockito.anyList(),
Mockito.anyList()))
+ .thenReturn(
+ CompletableFuture.completedFuture(Arrays.asList(StatusCode.GOOD,
StatusCode.GOOD)));
+ final IoTDBOpcUaClient client = createClient(miloClient);
+ final Map<IDeviceID, List<Pair<IMeasurementSchema, TimeValuePair>>>
deviceLastValues =
+ new LinkedHashMap<>();
+ deviceLastValues.put(
+ IDeviceID.Factory.DEFAULT_FACTORY.create("root.db.d1"),
+ Collections.singletonList(lastValue("s1", 1L, 11L)));
+ deviceLastValues.put(
+ IDeviceID.Factory.DEFAULT_FACTORY.create("root.db.d2"),
+ Collections.singletonList(lastValue("s1", 2L, 22L)));
+
+ client.transferLastValues(deviceLastValues, false, createSink());
+
+ Mockito.verify(miloClient)
+ .writeValuesAsync(
+ Mockito.argThat(nodeIds("root/db/d1/s1", "root/db/d2/s1")),
+ Mockito.argThat(listWithSize(2)));
+ }
+
@Test
public void testTransferCreatesAndRetriesOnlyMissingNodes() throws Exception
{
final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
@@ -157,6 +188,13 @@ public class IoTDBOpcUaClientTest {
return tablet;
}
+ private static Pair<IMeasurementSchema, TimeValuePair> lastValue(
+ final String measurement, final long timestamp, final long value) {
+ return new Pair<>(
+ new MeasurementSchema(measurement, TSDataType.INT64),
+ new TimeValuePair(timestamp,
TsPrimitiveType.getByType(TSDataType.INT64, value)));
+ }
+
private static ArgumentMatcher<List<NodeId>> nodeIds(final String...
identifiers) {
return nodeIds -> {
if (nodeIds.size() != identifiers.length) {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpaceMetadataTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpaceMetadataTest.java
new file mode 100644
index 00000000000..a7410aea71a
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpaceMetadataTest.java
@@ -0,0 +1,167 @@
+/*
+ * 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.iotdb.db.pipe.sink.protocol.opcua.server;
+
+import org.apache.iotdb.db.pipe.sink.protocol.opcua.OpcUaSink;
+
+import org.apache.tsfile.common.constant.TsFileConstant;
+import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.file.metadata.StringArrayDeviceID;
+import org.apache.tsfile.read.TimeValuePair;
+import org.apache.tsfile.utils.Binary;
+import org.apache.tsfile.utils.DateUtils;
+import org.apache.tsfile.utils.Pair;
+import org.apache.tsfile.utils.TsPrimitiveType;
+import org.apache.tsfile.write.schema.IMeasurementSchema;
+import org.apache.tsfile.write.schema.MeasurementSchema;
+import org.eclipse.milo.opcua.stack.core.types.builtin.DateTime;
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+import java.time.LocalDate;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Date;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.stream.Collectors;
+
+public class OpcUaNameSpaceMetadataTest {
+
+ @Test
+ public void testTransferLastValuesForTreeModel() throws Exception {
+ final CapturedRow capturedRow = new CapturedRow();
+
+ OpcUaNameSpace.transferLastValues(
+ IDeviceID.Factory.DEFAULT_FACTORY.create("root.sg.d1"),
+ Arrays.asList(
+ lastValue("s1", TSDataType.INT64, 5L, 50L),
+ lastValue("s2", TSDataType.TEXT, 6L, "last"),
+ new Pair<>(new MeasurementSchema("empty", TSDataType.INT64), null),
+ new Pair<>(
+ new MeasurementSchema(TsFileConstant.TIME_COLUMN_ID,
TSDataType.INT64),
+ timeValue(TSDataType.INT64, 7L, 8L))),
+ false,
+ createSink(),
+ capturedRow::capture);
+
+ Assert.assertArrayEquals(new String[] {"root", "sg", "d1"},
capturedRow.segments.get());
+ Assert.assertEquals(Arrays.asList("s1", "s2"),
capturedRow.getMeasurementNames());
+ Assert.assertEquals(Arrays.asList(5L, 6L), capturedRow.timestamps.get());
+ Assert.assertEquals(Arrays.asList(50L, "last"), capturedRow.values.get());
+ }
+
+ @Test
+ public void testTransferLastValuesForTableModel() throws Exception {
+ final int lastDate = DateUtils.parseDateExpressionToInt(LocalDate.of(2024,
1, 2));
+ final long lastTimestamp = 1_700_000_001_000L;
+ final CapturedRow capturedRow = new CapturedRow();
+ final OpcUaSink sink = createSink();
+ Mockito.when(sink.getDatabaseName()).thenReturn("database");
+ Mockito.when(sink.getPlaceHolder4NullTag()).thenReturn("null_tag");
+
+ OpcUaNameSpace.transferLastValues(
+ new StringArrayDeviceID("table", "tag", null, "tag2"),
+ Arrays.asList(
+ lastValue("date", TSDataType.DATE, 2L, lastDate),
+ lastValue("timestamp", TSDataType.TIMESTAMP, 4L, lastTimestamp)),
+ true,
+ sink,
+ capturedRow::capture);
+
+ Assert.assertArrayEquals(
+ new String[] {"database", "table", "tag", "null_tag", "tag2"},
capturedRow.segments.get());
+ Assert.assertEquals(Arrays.asList("date", "timestamp"),
capturedRow.getMeasurementNames());
+ Assert.assertEquals(Arrays.asList(2L, 4L), capturedRow.timestamps.get());
+ Assert.assertEquals(
+ new DateTime(new
Date(DateUtils.parseIntToDate(lastDate).getTime())).getUtcTime(),
+ ((DateTime) capturedRow.values.get().get(0)).getUtcTime());
+ Assert.assertEquals(
+ OpcUaNameSpace.timestampToUtc(lastTimestamp),
+ ((DateTime) capturedRow.values.get().get(1)).getUtcTime());
+ }
+
+ @Test
+ public void testTransferLastValuesSupportsBinaryValues() throws Exception {
+ final CapturedRow capturedRow = new CapturedRow();
+
+ OpcUaNameSpace.transferLastValues(
+ IDeviceID.Factory.DEFAULT_FACTORY.create("root.sg.d1"),
+ Arrays.asList(lastValue("blob", TSDataType.BLOB, 1L, "payload")),
+ false,
+ createSink(),
+ capturedRow::capture);
+
+ Assert.assertEquals(Arrays.asList("blob"),
capturedRow.getMeasurementNames());
+ Assert.assertEquals(TSDataType.BLOB,
capturedRow.schemas.get().get(0).getType());
+ Assert.assertEquals(Arrays.asList("payload"), capturedRow.values.get());
+ }
+
+ private static OpcUaSink createSink() {
+ final OpcUaSink sink = Mockito.mock(OpcUaSink.class);
+ Mockito.when(sink.getPlaceHolder4NullTag()).thenReturn("null");
+ return sink;
+ }
+
+ private static Pair<IMeasurementSchema, TimeValuePair> lastValue(
+ final String measurement,
+ final TSDataType dataType,
+ final long timestamp,
+ final Object value) {
+ return new Pair<>(
+ new MeasurementSchema(measurement, dataType), timeValue(dataType,
timestamp, value));
+ }
+
+ private static TimeValuePair timeValue(
+ final TSDataType dataType, final long timestamp, final Object value) {
+ final Object primitiveValue =
+ dataType == TSDataType.TEXT || dataType == TSDataType.BLOB || dataType
== TSDataType.STRING
+ ? new Binary(
+ String.valueOf(value),
org.apache.tsfile.common.conf.TSFileConfig.STRING_CHARSET)
+ : value;
+ return new TimeValuePair(timestamp, TsPrimitiveType.getByType(dataType,
primitiveValue));
+ }
+
+ private static class CapturedRow {
+ private final AtomicReference<String[]> segments = new AtomicReference<>();
+ private final AtomicReference<List<IMeasurementSchema>> schemas = new
AtomicReference<>();
+ private final AtomicReference<List<Long>> timestamps = new
AtomicReference<>();
+ private final AtomicReference<List<Object>> values = new
AtomicReference<>();
+
+ private void capture(
+ final String[] segments,
+ final List<IMeasurementSchema> schemas,
+ final List<Long> timestamps,
+ final List<Object> values,
+ final OpcUaSink sink) {
+ this.segments.set(segments);
+ this.schemas.set(new ArrayList<>(schemas));
+ this.timestamps.set(new ArrayList<>(timestamps));
+ this.values.set(new ArrayList<>(values));
+ }
+
+ private List<String> getMeasurementNames() {
+ return schemas.get().stream()
+ .map(IMeasurementSchema::getMeasurementName)
+ .collect(Collectors.toList());
+ }
+ }
+}