This is an automated email from the ASF dual-hosted git repository.

jt2594838 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 91a44fd8c00 [Pipe] Split OPC UA sink requests by advertised server 
limits (#18486)
91a44fd8c00 is described below

commit 91a44fd8c009a5dfa336953e4fd5aea04ee81586
Author: Caideyipi <[email protected]>
AuthorDate: Mon Aug 24 12:17:35 2026 +0800

    [Pipe] Split OPC UA sink requests by advertised server limits (#18486)
    
    * Pipe: respect OPC UA operation limits
    
    * Pipe: remove OPC UA server operation limits
    
    * Revert "Pipe: remove OPC UA server operation limits"
    
    This reverts commit b67c2656640bc7f9ec237c7ea5eafbdae6da9d63.
    
    * Localize OPC UA operation limit logs
---
 .../apache/iotdb/db/i18n/DataNodePipeMessages.java |   6 ++
 .../apache/iotdb/db/i18n/DataNodePipeMessages.java |   6 ++
 .../protocol/opcua/client/IoTDBOpcUaClient.java    | 113 ++++++++++++++++++---
 .../opcua/client/IoTDBOpcUaClientTest.java         |  94 +++++++++++++++++
 4 files changed, 203 insertions(+), 16 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
index 7da8a596579..4df5d7b4147 100644
--- 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
@@ -2590,4 +2590,10 @@ public final class DataNodePipeMessages {
       "Failed to release TsFile parser memory for Pipe {} (creation time {}) 
in DataRegion {} because no reservation exists.";
   public static final String 
LOG_PIPE_PROCESSOR_WORKER_ARG_HAS_BEEN_PROCESSING_THE_SAME_EVENT_FOR_ARG_MS_PIPE_ARG_DATAREGION_ARG_SUBTASK_ARG_EVENT_ARG_THREAD_STATE_ARG_STACK_ARG_63B40775
 =
       "Pipe processor worker {} has been processing the same event for {} ms. 
Pipe: {}, DataRegion: {}, subtask: {}, event: {}, thread state: {}. Stack:{}";
+  public static final String 
LOG_OPC_UA_SERVER_OPERATION_LIMITS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_5D2BCC90
 =
+      "OPC UA server operation limits: maxNodesPerWrite={}, 
maxNodesPerNodeManagement={}";
+  public static final String 
LOG_INTERRUPTED_WHILE_READING_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_357D46A4
 =
+      "Interrupted while reading OPC UA server operation limits, use defaults: 
maxNodesPerWrite={}, maxNodesPerNodeManagement={}";
+  public static final String 
LOG_FAILED_TO_READ_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_65460871
 =
+      "Failed to read OPC UA server operation limits, use defaults: 
maxNodesPerWrite={}, maxNodesPerNodeManagement={}";
 }
diff --git 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
index 35c88b9abfe..027e6b4c34c 100644
--- 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
@@ -2418,4 +2418,10 @@ public final class DataNodePipeMessages {
       "无法释放 Pipe {}(创建时间 {})在 DataRegion {} 中的 TsFile 解析器内存,因为不存在对应的预留。";
   public static final String 
LOG_PIPE_PROCESSOR_WORKER_ARG_HAS_BEEN_PROCESSING_THE_SAME_EVENT_FOR_ARG_MS_PIPE_ARG_DATAREGION_ARG_SUBTASK_ARG_EVENT_ARG_THREAD_STATE_ARG_STACK_ARG_63B40775
 =
       "Pipe processor worker {} 已连续处理同一 event {} 
ms。Pipe:{},DataRegion:{},subtask:{},event:{},线程状态:{}。栈:{}";
+  public static final String 
LOG_OPC_UA_SERVER_OPERATION_LIMITS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_5D2BCC90
 =
+      "OPC UA 服务器操作限制:maxNodesPerWrite={},maxNodesPerNodeManagement={}";
+  public static final String 
LOG_INTERRUPTED_WHILE_READING_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_357D46A4
 =
+      "读取 OPC UA 
服务器操作限制时被中断,使用默认值:maxNodesPerWrite={},maxNodesPerNodeManagement={}";
+  public static final String 
LOG_FAILED_TO_READ_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_65460871
 =
+      "读取 OPC UA 
服务器操作限制失败,使用默认值:maxNodesPerWrite={},maxNodesPerNodeManagement={}";
 }
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 83fc94e2ea2..f00ec845047 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
@@ -64,6 +64,7 @@ import javax.annotation.Nullable;
 
 import java.nio.file.Paths;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
@@ -74,10 +75,14 @@ import java.util.concurrent.ExecutionException;
 import static 
org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace.convertToOpcDataType;
 import static 
org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace.timestampToUtc;
 import static org.eclipse.milo.opcua.stack.core.StatusCodes.Bad_Timeout;
+import static 
org.eclipse.milo.opcua.stack.core.types.enumerated.TimestampsToReturn.Neither;
 
 public class IoTDBOpcUaClient {
   private static final Logger LOGGER = 
LoggerFactory.getLogger(OpcUaNameSpace.class);
 
+  private static final int DEFAULT_MAX_NODES_PER_WRITE = 10_000;
+  private static final int DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT = 250;
+
   // Customized nodes
   private static final int NAME_SPACE_INDEX = 2;
 
@@ -90,6 +95,8 @@ public class IoTDBOpcUaClient {
   private OpcUaClient client;
   private final boolean historizing;
   private ClientRunner runner;
+  private int maxNodesPerWrite = DEFAULT_MAX_NODES_PER_WRITE;
+  private int maxNodesPerNodeManagement = 
DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT;
 
   public IoTDBOpcUaClient(
       final String nodeUrl,
@@ -119,6 +126,63 @@ public class IoTDBOpcUaClient {
       }
       break;
     }
+    updateOperationLimits();
+  }
+
+  private void updateOperationLimits() {
+    try {
+      final List<DataValue> operationLimits =
+          client
+              .readValuesAsync(
+                  0.0,
+                  Neither,
+                  Arrays.asList(
+                      
Identifiers.Server_ServerCapabilities_OperationLimits_MaxNodesPerWrite,
+                      Identifiers
+                          
.Server_ServerCapabilities_OperationLimits_MaxNodesPerNodeManagement))
+              .get();
+      maxNodesPerWrite = getOperationLimit(operationLimits, 0, 
DEFAULT_MAX_NODES_PER_WRITE);
+      maxNodesPerNodeManagement =
+          getOperationLimit(operationLimits, 1, 
DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT);
+      LOGGER.info(
+          DataNodePipeMessages
+              
.LOG_OPC_UA_SERVER_OPERATION_LIMITS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_5D2BCC90,
+          maxNodesPerWrite,
+          maxNodesPerNodeManagement);
+    } catch (final InterruptedException e) {
+      Thread.currentThread().interrupt();
+      LOGGER.warn(
+          DataNodePipeMessages
+              
.LOG_INTERRUPTED_WHILE_READING_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_357D46A4,
+          DEFAULT_MAX_NODES_PER_WRITE,
+          DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT);
+    } catch (final Exception e) {
+      LOGGER.warn(
+          DataNodePipeMessages
+              
.LOG_FAILED_TO_READ_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_65460871,
+          DEFAULT_MAX_NODES_PER_WRITE,
+          DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT,
+          e);
+    }
+  }
+
+  private static int getOperationLimit(
+      final List<DataValue> operationLimits, final int index, final int 
defaultValue) {
+    if (Objects.isNull(operationLimits) || operationLimits.size() <= index) {
+      return defaultValue;
+    }
+
+    final DataValue dataValue = operationLimits.get(index);
+    if (Objects.isNull(dataValue)
+        || Objects.isNull(dataValue.getStatusCode())
+        || !dataValue.getStatusCode().isGood()
+        || Objects.isNull(dataValue.getValue())
+        || !(dataValue.getValue().getValue() instanceof Number)) {
+      return defaultValue;
+    }
+
+    final long limit = ((Number) dataValue.getValue().getValue()).longValue();
+    return limit == 0 ? Integer.MAX_VALUE : (int) Math.min(limit, 
Integer.MAX_VALUE);
   }
 
   // Only support tree model & client-server
@@ -267,17 +331,23 @@ public class IoTDBOpcUaClient {
       }
     }
 
-    final AddNodesResponse addStatus = client.addNodesAsync(nodesToAdd).get();
-    for (final AddNodesResult result : addStatus.getResults()) {
-      if (!result.getStatusCode().equals(StatusCode.GOOD)
-          && result.getStatusCode().getValue() != 
StatusCodes.Bad_NodeIdExists) {
-        throw new PipeException(
-            DataNodePipeMessages.FAILED_TO_CREATE_NODES_AFTER_TRANSFER_DATA
-                + addStatus
-                + writeRequests
-                    .get(0)
-                    .getErrorString(new 
StatusCode(StatusCodes.Bad_NodeIdUnknown)));
+    for (int startIndex = 0; startIndex < nodesToAdd.size(); ) {
+      final int endIndex =
+          getBatchEndIndex(startIndex, nodesToAdd.size(), 
maxNodesPerNodeManagement);
+      final AddNodesResponse addStatus =
+          client.addNodesAsync(nodesToAdd.subList(startIndex, endIndex)).get();
+      for (final AddNodesResult result : addStatus.getResults()) {
+        if (!result.getStatusCode().equals(StatusCode.GOOD)
+            && result.getStatusCode().getValue() != 
StatusCodes.Bad_NodeIdExists) {
+          throw new PipeException(
+              DataNodePipeMessages.FAILED_TO_CREATE_NODES_AFTER_TRANSFER_DATA
+                  + addStatus
+                  + writeRequests
+                      .get(0)
+                      .getErrorString(new 
StatusCode(StatusCodes.Bad_NodeIdUnknown)));
+        }
       }
+      startIndex = endIndex;
     }
   }
 
@@ -294,13 +364,24 @@ public class IoTDBOpcUaClient {
 
   private List<StatusCode> writeValuesOnce(final List<OpcUaWriteRequest> 
writeRequests)
       throws Exception {
-    final List<NodeId> nodeIds = new ArrayList<>(writeRequests.size());
-    final List<DataValue> dataValues = new ArrayList<>(writeRequests.size());
-    for (final OpcUaWriteRequest writeRequest : writeRequests) {
-      nodeIds.add(writeRequest.nodeId);
-      dataValues.add(writeRequest.dataValue);
+    final List<StatusCode> writeStatuses = new 
ArrayList<>(writeRequests.size());
+    for (int startIndex = 0; startIndex < writeRequests.size(); ) {
+      final int endIndex = getBatchEndIndex(startIndex, writeRequests.size(), 
maxNodesPerWrite);
+      final List<NodeId> nodeIds = new ArrayList<>(endIndex - startIndex);
+      final List<DataValue> dataValues = new ArrayList<>(endIndex - 
startIndex);
+      for (int i = startIndex; i < endIndex; ++i) {
+        nodeIds.add(writeRequests.get(i).nodeId);
+        dataValues.add(writeRequests.get(i).dataValue);
+      }
+      writeStatuses.addAll(client.writeValuesAsync(nodeIds, dataValues).get());
+      startIndex = endIndex;
     }
-    return client.writeValuesAsync(nodeIds, dataValues).get();
+    return writeStatuses;
+  }
+
+  private static int getBatchEndIndex(
+      final int startIndex, final int totalSize, final int batchSize) {
+    return (int) Math.min((long) totalSize, (long) startIndex + batchSize);
   }
 
   private static final class OpcUaWriteRequest {
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 8f8333143e3..d327155f32d 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
@@ -34,8 +34,12 @@ import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
 import org.eclipse.milo.opcua.sdk.client.identity.AnonymousProvider;
 import org.eclipse.milo.opcua.stack.core.StatusCodes;
 import org.eclipse.milo.opcua.stack.core.security.SecurityPolicy;
+import org.eclipse.milo.opcua.stack.core.types.builtin.DataValue;
+import org.eclipse.milo.opcua.stack.core.types.builtin.ExpandedNodeId;
 import org.eclipse.milo.opcua.stack.core.types.builtin.NodeId;
 import org.eclipse.milo.opcua.stack.core.types.builtin.StatusCode;
+import org.eclipse.milo.opcua.stack.core.types.builtin.Variant;
+import org.eclipse.milo.opcua.stack.core.types.enumerated.TimestampsToReturn;
 import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesItem;
 import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesResponse;
 import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesResult;
@@ -52,6 +56,8 @@ import java.util.List;
 import java.util.Map;
 import java.util.concurrent.CompletableFuture;
 
+import static 
org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.uint;
+
 public class IoTDBOpcUaClientTest {
 
   @Test
@@ -70,6 +76,30 @@ public class IoTDBOpcUaClientTest {
             Mockito.argThat(listWithSize(2)));
   }
 
+  @Test
+  public void testTransferSplitsWritesAtServerLimit() throws Exception {
+    final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
+    Mockito.when(miloClient.writeValuesAsync(Mockito.anyList(), 
Mockito.anyList()))
+        .thenAnswer(
+            invocation -> {
+              final int size = ((List<?>) invocation.getArguments()[0]).size();
+              return 
CompletableFuture.completedFuture(Collections.nCopies(size, StatusCode.GOOD));
+            });
+    final IoTDBOpcUaClient client = createClient(miloClient, 1, 250);
+
+    client.transfer(createTablet(), createSink());
+
+    final InOrder inOrder = Mockito.inOrder(miloClient);
+    inOrder
+        .verify(miloClient)
+        .writeValuesAsync(
+            Mockito.argThat(nodeIds("root/db/d1/s1")), 
Mockito.argThat(listWithSize(1)));
+    inOrder
+        .verify(miloClient)
+        .writeValuesAsync(
+            Mockito.argThat(nodeIds("root/db/d1/s2")), 
Mockito.argThat(listWithSize(1)));
+  }
+
   @Test
   public void testTransferLastValuesBatchesDevicesInOneRequest() throws 
Exception {
     final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
@@ -133,6 +163,56 @@ public class IoTDBOpcUaClientTest {
             Mockito.argThat(nodeIds("root/db/d1/s1")), 
Mockito.argThat(listWithSize(1)));
   }
 
+  @Test
+  public void testTransferSplitsMissingNodeCreationAtServerLimit() throws 
Exception {
+    final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
+    Mockito.when(miloClient.writeValuesAsync(Mockito.anyList(), 
Mockito.anyList()))
+        .thenReturn(
+            CompletableFuture.completedFuture(
+                Arrays.asList(
+                    new StatusCode(StatusCodes.Bad_NodeIdUnknown),
+                    new StatusCode(StatusCodes.Bad_NodeIdUnknown))))
+        .thenReturn(
+            CompletableFuture.completedFuture(Arrays.asList(StatusCode.GOOD, 
StatusCode.GOOD)));
+
+    final AddNodesResponse addNodesResponse = 
Mockito.mock(AddNodesResponse.class);
+    final AddNodesResult addNodesResult = Mockito.mock(AddNodesResult.class);
+    Mockito.when(addNodesResult.getStatusCode()).thenReturn(StatusCode.GOOD);
+    Mockito.when(addNodesResponse.getResults()).thenReturn(new 
AddNodesResult[] {addNodesResult});
+    Mockito.when(miloClient.addNodesAsync(Mockito.anyList()))
+        .thenReturn(CompletableFuture.completedFuture(addNodesResponse));
+
+    final IoTDBOpcUaClient client = Mockito.spy(createClient(miloClient, 
10_000, 1));
+    final AddNodesItem firstNode = Mockito.mock(AddNodesItem.class);
+    final AddNodesItem secondNode = Mockito.mock(AddNodesItem.class);
+    final ExpandedNodeId firstNodeId = new NodeId(2, 
"root/db/d1/s1").expanded();
+    final ExpandedNodeId secondNodeId = new NodeId(2, 
"root/db/d1/s2").expanded();
+    Mockito.when(firstNode.getRequestedNewNodeId()).thenReturn(firstNodeId);
+    Mockito.when(secondNode.getRequestedNewNodeId()).thenReturn(secondNodeId);
+    Mockito.doAnswer(
+            invocation ->
+                Collections.singletonList(
+                    "s1".equals(invocation.getArguments()[1]) ? firstNode : 
secondNode))
+        .when(client)
+        .getNodesToAdd(
+            Mockito.any(String[].class),
+            Mockito.anyString(),
+            Mockito.any(NodeId.class),
+            Mockito.any());
+
+    client.transfer(createTablet(), createSink());
+
+    Mockito.verify(miloClient, 
Mockito.times(2)).addNodesAsync(Mockito.argThat(listWithSize(1)));
+    final InOrder inOrder = Mockito.inOrder(miloClient);
+    inOrder
+        .verify(miloClient)
+        .writeValuesAsync(Mockito.argThat(listWithSize(2)), 
Mockito.argThat(listWithSize(2)));
+    inOrder.verify(miloClient, 
Mockito.times(2)).addNodesAsync(Mockito.argThat(listWithSize(1)));
+    inOrder
+        .verify(miloClient)
+        .writeValuesAsync(Mockito.argThat(listWithSize(2)), 
Mockito.argThat(listWithSize(2)));
+  }
+
   @Test
   public void testTransferFailsOnNonRecoverableStatus() throws Exception {
     final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
@@ -154,6 +234,12 @@ public class IoTDBOpcUaClientTest {
   }
 
   private static IoTDBOpcUaClient createClient(final OpcUaClient miloClient) 
throws Exception {
+    return createClient(miloClient, 10_000, 250);
+  }
+
+  private static IoTDBOpcUaClient createClient(
+      final OpcUaClient miloClient, final int maxNodesPerWrite, final int 
maxNodesPerNodeManagement)
+      throws Exception {
     final IoTDBOpcUaClient client =
         new IoTDBOpcUaClient(
             "opc.tcp://127.0.0.1:12686", SecurityPolicy.None, new 
AnonymousProvider(), false);
@@ -163,6 +249,14 @@ public class IoTDBOpcUaClientTest {
     final CompletableFuture<OpcUaClient> connectFuture =
         CompletableFuture.completedFuture(miloClient);
     Mockito.when(miloClient.connectAsync()).thenReturn(connectFuture);
+    Mockito.when(
+            miloClient.readValuesAsync(
+                Mockito.anyDouble(), Mockito.eq(TimestampsToReturn.Neither), 
Mockito.anyList()))
+        .thenReturn(
+            CompletableFuture.completedFuture(
+                Arrays.asList(
+                    new DataValue(new Variant(uint(maxNodesPerWrite))),
+                    new DataValue(new 
Variant(uint(maxNodesPerNodeManagement))))));
     client.run(miloClient);
     return client;
   }

Reply via email to