jt2594838 commented on code in PR #18433:
URL: https://github.com/apache/iotdb/pull/18433#discussion_r3748619087


##########
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java:
##########
@@ -130,6 +134,98 @@ public void transfer(final Tablet tablet, final boolean 
isTableModel, final OpcU
     }
   }
 
+  /**
+   * 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();
+    switch (dataType) {

Review Comment:
   Use enhanced-switch



##########
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java:
##########
@@ -130,6 +134,98 @@ public void transfer(final Tablet tablet, final boolean 
isTableModel, final OpcU
     }
   }
 
+  /**
+   * 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();
+    switch (dataType) {
+      case DATE:
+        return new DateTime(
+            new java.util.Date(DateUtils.parseIntToDate(((Number) 
value).intValue()).getTime()));

Review Comment:
   Mind the import, add to your local skill



##########
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java:
##########
@@ -497,6 +512,151 @@ public void heartbeat() throws Exception {
     // 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;
+      }
+
+      final boolean isTableModel = 
pipeTsFileInsertionEvent.isTableModelEvent();
+      for (final Map.Entry<IDeviceID, List<Pair<IMeasurementSchema, 
TimeValuePair>>> entry :
+          deviceLastValues.entrySet()) {
+        if (Objects.nonNull(nameSpace)) {
+          nameSpace.transferLastValues(entry.getKey(), entry.getValue(), 
isTableModel, this);
+        } else if (Objects.nonNull(client)) {
+          client.transferLastValues(entry.getKey(), entry.getValue(), 
isTableModel, this);
+        } else {
+          throw new 
PipeException(DataNodePipeMessages.NO_OPC_CLIENT_OR_SERVER_IS_SPECIFIED);
+        }
+      }

Review Comment:
   Will the devices be batched?
   For many-device-few-measurement scenarios, transferring values one by one 
may be slow.



##########
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java:
##########
@@ -497,6 +512,151 @@ public void heartbeat() throws Exception {
     // 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;
+      }
+
+      final boolean isTableModel = 
pipeTsFileInsertionEvent.isTableModelEvent();
+      for (final Map.Entry<IDeviceID, List<Pair<IMeasurementSchema, 
TimeValuePair>>> entry :
+          deviceLastValues.entrySet()) {
+        if (Objects.nonNull(nameSpace)) {
+          nameSpace.transferLastValues(entry.getKey(), entry.getValue(), 
isTableModel, this);
+        } else if (Objects.nonNull(client)) {
+          client.transferLastValues(entry.getKey(), entry.getValue(), 
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 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)) {
+          throw new IOException();
+        }

Review Comment:
   Add clearer message. Or return null instead.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to