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;
}