This is an automated email from the ASF dual-hosted git repository.
Caideyipi pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/dev/1.3 by this push:
new f9723aeb7a1 [Pipe] Optimize OPC UA TsFile last-value transfer (#18433)
(#18491)
f9723aeb7a1 is described below
commit f9723aeb7a108a0702925150847ae421d2496df7
Author: Caideyipi <[email protected]>
AuthorDate: Fri Aug 21 09:33:46 2026 +0800
[Pipe] Optimize OPC UA TsFile last-value transfer (#18433) (#18491)
(cherry picked from commit dc2e7c861b62b049b1b6cc26b9eb2fd7d3245a0d)
---
.../db/pipe/sink/protocol/opcua/OpcUaSink.java | 161 +++++++++++++++++++++
.../protocol/opcua/client/IoTDBOpcUaClient.java | 37 ++++-
.../sink/protocol/opcua/server/OpcUaNameSpace.java | 84 ++++++++++-
.../opcua/OpcUaSinkTsFileMetadataTest.java | 105 ++++++++++++++
.../opcua/client/IoTDBOpcUaClientTest.java | 36 +++++
.../opcua/server/OpcUaNameSpaceMetadataTest.java | 128 ++++++++++++++++
6 files changed, 546 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 d6cf17c07ad..fa77611c12a 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
@@ -22,6 +22,7 @@ package org.apache.iotdb.db.pipe.sink.protocol.opcua;
import org.apache.iotdb.db.conf.IoTDBConfig;
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;
@@ -32,11 +33,20 @@ 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.MeasurementSchema;
import org.eclipse.milo.opcua.sdk.client.api.identity.AnonymousProvider;
import org.eclipse.milo.opcua.sdk.client.api.identity.IdentityProvider;
import org.eclipse.milo.opcua.sdk.client.api.identity.UsernameProvider;
@@ -49,7 +59,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;
@@ -467,6 +481,153 @@ 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.shouldParseTime();
+ }
+
+ 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<MeasurementSchema, 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;
+ }
+
+ if (Objects.nonNull(nameSpace)) {
+ for (final Map.Entry<IDeviceID, List<Pair<MeasurementSchema,
TimeValuePair>>> entry :
+ deviceLastValues.entrySet()) {
+ nameSpace.transferLastValues(entry.getKey(), entry.getValue(), 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, this);
+ } else {
+ throw new PipeException("No OPC client or server is specified when
transferring TsFile");
+ }
+ return TsFileTransferResult.TRANSFERRED;
+ } finally {
+
pipeTsFileInsertionEvent.decreaseReferenceCount(OpcUaSink.class.getName(),
false);
+ }
+ }
+
+ static @Nullable Map<IDeviceID, List<Pair<MeasurementSchema,
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<MeasurementSchema, 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<MeasurementSchema, 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 cf8ab9001fb..65c651018ef 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
@@ -27,6 +27,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.MeasurementSchema;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
@@ -62,6 +65,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;
@@ -122,6 +126,24 @@ public class IoTDBOpcUaClient {
tablet, sink, this::transferTabletRowForClientServerModel);
}
+ public void transferLastValues(
+ final Map<IDeviceID, List<Pair<MeasurementSchema, TimeValuePair>>>
deviceLastValues,
+ final OpcUaSink sink)
+ throws Exception {
+ final List<OpcUaWriteRequest> writeRequests = new ArrayList<>();
+ for (final Map.Entry<IDeviceID, List<Pair<MeasurementSchema,
TimeValuePair>>> entry :
+ deviceLastValues.entrySet()) {
+ OpcUaNameSpace.transferLastValues(
+ entry.getKey(),
+ entry.getValue(),
+ sink,
+ (segments, measurementSchemas, timestamps, values, currentSink) ->
+ collectWriteRequests(
+ segments, measurementSchemas, timestamps, values,
currentSink, writeRequests));
+ }
+ writeValues(writeRequests);
+ }
+
private void transferTabletRowForClientServerModel(
final String[] segments,
final List<MeasurementSchema> measurementSchemas,
@@ -129,11 +151,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<MeasurementSchema> 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))) {
@@ -176,8 +209,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 64fd9fce310..c61b0ca5d97 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
@@ -30,7 +30,12 @@ import org.apache.iotdb.pipe.api.event.Event;
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.PlainDeviceID;
+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.MeasurementSchema;
@@ -60,11 +65,11 @@ 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.Date;
import java.util.List;
import java.util.Objects;
import java.util.Set;
@@ -123,6 +128,79 @@ 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<MeasurementSchema, TimeValuePair>> lastValues,
+ final OpcUaSink sink)
+ throws Exception {
+ transferLastValues(deviceID, lastValues, sink,
this::transferTabletRowForClientServerModel);
+ }
+
+ public static void transferLastValues(
+ final IDeviceID deviceID,
+ final List<Pair<MeasurementSchema, TimeValuePair>> lastValues,
+ final OpcUaSink sink,
+ final TabletRowConsumer consumer)
+ throws Exception {
+ // The 1.3 branch only supports tree-model device IDs.
+ final String[] segments = ((PlainDeviceID)
deviceID).toStringID().split("\\.");
+
+ final List<MeasurementSchema> schemas = new ArrayList<>(lastValues.size());
+ final List<Long> timestamps = new ArrayList<>(lastValues.size());
+ final List<Object> values = new ArrayList<>(lastValues.size());
+ for (final Pair<MeasurementSchema, TimeValuePair> lastValue : lastValues) {
+ if (Objects.isNull(lastValue)
+ || Objects.isNull(lastValue.getLeft())
+ || Objects.isNull(lastValue.getLeft().getMeasurementId())
+ ||
TsFileConstant.TIME_COLUMN_ID.equals(lastValue.getLeft().getMeasurementId())
+ || 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();
+ switch (dataType) {
+ case DATE:
+ return new DateTime(
+ new Date(DateUtils.parseIntToDate(((Number)
value).intValue()).getTime()));
+ case TIMESTAMP:
+ return new DateTime(timestampToUtc(((Number) value).longValue()));
+ case TEXT:
+ case BLOB:
+ case STRING:
+ return value instanceof Binary ? value.toString() :
String.valueOf(value);
+ case BOOLEAN:
+ case INT32:
+ case INT64:
+ case FLOAT:
+ case DOUBLE:
+ return value;
+ case VECTOR:
+ case UNKNOWN:
+ default:
+ throw new UnSupportedDataTypeException("UnSupported dataType " +
dataType);
+ }
+ }
+
public static void transferTabletForClientServerModel(
final Tablet tablet, final OpcUaSink sink, final TabletRowConsumer
consumer)
throws Exception {
@@ -323,7 +401,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..a9e9da3fc55
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSinkTsFileMetadataTest.java
@@ -0,0 +1,105 @@
+/*
+ * 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.TSDataType;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.read.TimeValuePair;
+import org.apache.tsfile.read.common.Path;
+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.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<MeasurementSchema> 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));
+ tablet.rowSize = 3;
+
+ try (final TsFileWriter writer = new TsFileWriter(tsFile)) {
+ writer.registerTimeseries(new Path(device), schemas);
+ writer.write(tablet);
+ }
+
+ final Map<IDeviceID, List<Pair<MeasurementSchema, 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");
+ }
+
+ private static Map<String, TimeValuePair> toMeasurementMap(
+ final List<Pair<MeasurementSchema, TimeValuePair>> lastValues) {
+ final Map<String, TimeValuePair> result = new LinkedHashMap<>();
+ lastValues.forEach(
+ lastValue -> result.put(lastValue.getLeft().getMeasurementId(),
lastValue.getRight()));
+ return result;
+ }
+
+ 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/client/IoTDBOpcUaClientTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java
index 98f8cfaa88d..0eaa1cf0234 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,6 +23,11 @@ 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.file.metadata.PlainDeviceID;
+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.MeasurementSchema;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
@@ -43,7 +48,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 {
@@ -64,6 +71,28 @@ public class IoTDBOpcUaClientTest {
Mockito.argThat(listWithSize(2)));
}
+ @Test
+ public void testTransferLastValuesBatchesDevicesInOneRequest() throws
Exception {
+ final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
+ Mockito.when(miloClient.writeValues(Mockito.anyList(), Mockito.anyList()))
+ .thenReturn(
+ CompletableFuture.completedFuture(Arrays.asList(StatusCode.GOOD,
StatusCode.GOOD)));
+ final IoTDBOpcUaClient client = createClient(miloClient);
+ final Map<IDeviceID, List<Pair<MeasurementSchema, TimeValuePair>>>
deviceLastValues =
+ new LinkedHashMap<>();
+ deviceLastValues.put(
+ new PlainDeviceID("root.db.d1"),
Collections.singletonList(lastValue("s1", 1L, 11L)));
+ deviceLastValues.put(
+ new PlainDeviceID("root.db.d2"),
Collections.singletonList(lastValue("s1", 2L, 22L)));
+
+ client.transferLastValues(deviceLastValues, createSink());
+
+ Mockito.verify(miloClient)
+ .writeValues(
+ 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);
@@ -156,6 +185,13 @@ public class IoTDBOpcUaClientTest {
return tablet;
}
+ private static Pair<MeasurementSchema, 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..c9ce3ff3a78
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpaceMetadataTest.java
@@ -0,0 +1,128 @@
+/*
+ * 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.PlainDeviceID;
+import org.apache.tsfile.read.TimeValuePair;
+import org.apache.tsfile.utils.Binary;
+import org.apache.tsfile.utils.Pair;
+import org.apache.tsfile.utils.TsPrimitiveType;
+import org.apache.tsfile.write.schema.MeasurementSchema;
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+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(
+ new PlainDeviceID("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))),
+ 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 testTransferLastValuesSupportsBinaryValues() throws Exception {
+ final CapturedRow capturedRow = new CapturedRow();
+
+ OpcUaNameSpace.transferLastValues(
+ new PlainDeviceID("root.sg.d1"),
+ Arrays.asList(lastValue("blob", TSDataType.BLOB, 1L, "payload")),
+ 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() {
+ return Mockito.mock(OpcUaSink.class);
+ }
+
+ private static Pair<MeasurementSchema, 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<MeasurementSchema>> 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<MeasurementSchema> 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(MeasurementSchema::getMeasurementId)
+ .collect(Collectors.toList());
+ }
+ }
+}