From 86a5a7dc92409b94a930ed962de9094464831d59 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Mon, 20 Jul 2026 17:44:50 +0800 Subject: [PATCH] [Performance] Optimize tablet RPC deserialization (#18248) * [Performance] Optimize tablet RPC deserialization * [Performance] Add tablet RPC deserialization benchmarks * Address tablet deserialization review comments --- .../thrift/impl/ClientRPCServiceImpl.java | 6 +- .../plan/parser/StatementGenerator.java | 45 ++++++----- .../apache/iotdb/commons/utils/PathUtils.java | 68 +++++++++++++++- .../iotdb/commons/utils/PathUtilsTest.java | 77 +++++++++++++++++++ 4 files changed, 166 insertions(+), 30 deletions(-) create mode 100644 iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/utils/PathUtilsTest.java diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java index 3fe0e5e526789..c4bfd21a7a406 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java @@ -2058,9 +2058,7 @@ public TSStatus insertTablets(TSInsertTabletsReq req) { if (!SESSION_MANAGER.checkLogin(clientSession)) { return getNotLoggedInStatus(); } - - req.setMeasurementsList( - PathUtils.checkIsLegalSingleMeasurementListsAndUpdate(req.getMeasurementsList())); + PathUtils.checkIsLegalSingleMeasurementListsAndUpdateInPlace(req.getMeasurementsList()); // Step 1: transfer from TSInsertTabletsReq to Statement InsertMultiTabletsStatement statement = StatementGenerator.createStatement(req); @@ -2119,7 +2117,7 @@ public TSStatus insertTablet(TSInsertTabletReq req) { } // check whether measurement is legal according to syntax convention - req.setMeasurements(PathUtils.checkIsLegalSingleMeasurementsAndUpdate(req.getMeasurements())); + PathUtils.checkIsLegalSingleMeasurementsAndUpdateInPlace(req.getMeasurements()); // Step 1: transfer from TSInsertTabletReq to Statement InsertTabletStatement statement = StatementGenerator.createStatement(req); // return success when this statement is empty because server doesn't need to execute it diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/StatementGenerator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/StatementGenerator.java index 9f149abd531eb..9e19889ee2b46 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/StatementGenerator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/StatementGenerator.java @@ -313,6 +313,7 @@ public static InsertTabletStatement createStatement(TSInsertTabletReq insertTabl insertStatement.setDevicePath( DEVICE_PATH_CACHE.getPartialPath(insertTabletReq.getPrefixPath())); insertStatement.setMeasurements(insertTabletReq.getMeasurements().toArray(new String[0])); + TSDataType[] dataTypes = deserializeDataTypes(insertTabletReq.types); long[] timestamps = QueryDataSetUtils.readTimesFromBuffer(insertTabletReq.timestamps, insertTabletReq.size); if (timestamps.length != 0) { @@ -321,19 +322,12 @@ public static InsertTabletStatement createStatement(TSInsertTabletReq insertTabl insertStatement.setTimes(timestamps); insertStatement.setColumns( QueryDataSetUtils.readTabletValuesFromBuffer( - insertTabletReq.values, - insertTabletReq.types, - insertTabletReq.types.size(), - insertTabletReq.size)); + insertTabletReq.values, dataTypes, dataTypes.length, insertTabletReq.size)); insertStatement.setBitMaps( QueryDataSetUtils.readBitMapsFromBuffer( - insertTabletReq.values, insertTabletReq.types.size(), insertTabletReq.size) + insertTabletReq.values, dataTypes.length, insertTabletReq.size) .orElse(null)); insertStatement.setRowCount(insertTabletReq.size); - TSDataType[] dataTypes = new TSDataType[insertTabletReq.types.size()]; - for (int i = 0; i < insertTabletReq.types.size(); i++) { - dataTypes[i] = TSDataType.deserialize((byte) insertTabletReq.types.get(i).intValue()); - } insertStatement.setDataTypes(dataTypes); insertStatement.setAligned(insertTabletReq.isAligned); PERFORMANCE_OVERVIEW_METRICS.recordParseCost(System.nanoTime() - startTime); @@ -345,32 +339,29 @@ public static InsertMultiTabletsStatement createStatement(TSInsertTabletsReq req final long startTime = System.nanoTime(); // construct insert statement InsertMultiTabletsStatement insertStatement = new InsertMultiTabletsStatement(); - List insertTabletStatementList = new ArrayList<>(); - for (int i = 0; i < req.prefixPaths.size(); i++) { + int tabletCount = req.prefixPaths.size(); + List insertTabletStatementList = new ArrayList<>(tabletCount); + for (int i = 0; i < tabletCount; i++) { + List measurements = req.measurementsList.get(i); + TSDataType[] dataTypes = deserializeDataTypes(req.typesList.get(i)); + int rowCount = req.sizeList.get(i); InsertTabletStatement insertTabletStatement = new InsertTabletStatement(); insertTabletStatement.setDevicePath(DEVICE_PATH_CACHE.getPartialPath(req.prefixPaths.get(i))); - insertTabletStatement.setMeasurements(req.measurementsList.get(i).toArray(new String[0])); + insertTabletStatement.setMeasurements(measurements.toArray(new String[0])); long[] timestamps = - QueryDataSetUtils.readTimesFromBuffer(req.timestampsList.get(i), req.sizeList.get(i)); + QueryDataSetUtils.readTimesFromBuffer(req.timestampsList.get(i), rowCount); if (timestamps.length != 0) { TimestampPrecisionUtils.checkTimestampPrecision(timestamps[timestamps.length - 1]); } insertTabletStatement.setTimes(timestamps); insertTabletStatement.setColumns( QueryDataSetUtils.readTabletValuesFromBuffer( - req.valuesList.get(i), - req.typesList.get(i), - req.measurementsList.get(i).size(), - req.sizeList.get(i))); + req.valuesList.get(i), dataTypes, measurements.size(), rowCount)); insertTabletStatement.setBitMaps( QueryDataSetUtils.readBitMapsFromBuffer( - req.valuesList.get(i), req.measurementsList.get(i).size(), req.sizeList.get(i)) + req.valuesList.get(i), measurements.size(), rowCount) .orElse(null)); - insertTabletStatement.setRowCount(req.sizeList.get(i)); - TSDataType[] dataTypes = new TSDataType[req.typesList.get(i).size()]; - for (int j = 0; j < dataTypes.length; j++) { - dataTypes[j] = TSDataType.deserialize((byte) req.typesList.get(i).get(j).intValue()); - } + insertTabletStatement.setRowCount(rowCount); insertTabletStatement.setDataTypes(dataTypes); insertTabletStatement.setAligned(req.isAligned); // skip empty tablet @@ -384,6 +375,14 @@ public static InsertMultiTabletsStatement createStatement(TSInsertTabletsReq req return insertStatement; } + private static TSDataType[] deserializeDataTypes(List serializedDataTypes) { + TSDataType[] dataTypes = new TSDataType[serializedDataTypes.size()]; + for (int i = 0; i < dataTypes.length; i++) { + dataTypes[i] = TSDataType.deserialize((byte) serializedDataTypes.get(i).intValue()); + } + return dataTypes; + } + public static InsertRowsStatement createStatement(TSInsertRecordsReq req) throws IllegalPathException, QueryProcessException { final long startTime = System.nanoTime(); diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/PathUtils.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/PathUtils.java index e3af7aa5eb91a..41a77d3be3181 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/PathUtils.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/PathUtils.java @@ -69,7 +69,7 @@ public static List> checkIsLegalSingleMeasurementListsAndUpdate( } // skip checking duplicated measurements Map checkedMeasurements = new HashMap<>(); - List> res = new ArrayList<>(); + List> res = new ArrayList<>(measurementLists.size()); for (List measurements : measurementLists) { res.add(checkLegalSingleMeasurementsAndSkipDuplicate(measurements, checkedMeasurements)); } @@ -85,7 +85,7 @@ public static List checkLegalSingleMeasurementsAndSkipDuplicate( if (measurements == null) { return null; } - List res = new ArrayList<>(); + List res = new ArrayList<>(measurements.size()); for (String measurement : measurements) { if (measurement == null) { res.add(null); @@ -110,7 +110,7 @@ public static List checkIsLegalSingleMeasurementsAndUpdate(List if (measurements == null) { return null; } - List res = new ArrayList<>(); + List res = new ArrayList<>(measurements.size()); for (String measurement : measurements) { if (measurement == null) { continue; @@ -120,6 +120,68 @@ public static List checkIsLegalSingleMeasurementsAndUpdate(List return res; } + /** + * Check and canonicalize single measurements in place. This avoids allocating another list when + * the input is a mutable list created by Thrift. + */ + public static void checkIsLegalSingleMeasurementsAndUpdateInPlace(List measurements) + throws MetadataException { + if (measurements == null) { + return; + } + for (int i = 0; i < measurements.size(); ) { + String measurement = measurements.get(i); + if (measurement == null) { + measurements.remove(i); + } else { + measurements.set(i, checkAndReturnSingleMeasurement(measurement)); + i++; + } + } + } + + /** + * Check and canonicalize lists of single measurements in place. Duplicate measurements in one + * request are checked only once. + */ + public static void checkIsLegalSingleMeasurementListsAndUpdateInPlace( + List> measurementLists) throws MetadataException { + if (measurementLists == null || measurementLists.isEmpty()) { + return; + } + if (measurementLists.size() == 1) { + List measurements = measurementLists.get(0); + if (measurements == null) { + return; + } + for (int i = 0; i < measurements.size(); i++) { + String measurement = measurements.get(i); + if (measurement != null) { + measurements.set(i, checkAndReturnSingleMeasurement(measurement)); + } + } + return; + } + Map checkedMeasurements = new HashMap<>(); + for (List measurements : measurementLists) { + if (measurements == null) { + continue; + } + for (int i = 0; i < measurements.size(); i++) { + String measurement = measurements.get(i); + if (measurement == null) { + continue; + } + String checked = checkedMeasurements.get(measurement); + if (checked == null) { + checked = checkAndReturnSingleMeasurement(measurement); + checkedMeasurements.put(measurement, checked); + } + measurements.set(i, checked); + } + } + } + /** * check whether measurement is legal according to syntax convention measurement could be like a.b * (more than one node name), in template? diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/utils/PathUtilsTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/utils/PathUtilsTest.java new file mode 100644 index 0000000000000..1119e54737f7a --- /dev/null +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/utils/PathUtilsTest.java @@ -0,0 +1,77 @@ +/* + * 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.commons.utils; + +import org.apache.iotdb.commons.exception.MetadataException; + +import org.junit.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; + +public class PathUtilsTest { + + @Test + public void testCheckSingleMeasurementsInPlace() throws MetadataException { + List measurements = + new ArrayList<>(Arrays.asList("path_utils_test_s1", "`path_utils_test_s2`", null)); + + PathUtils.checkIsLegalSingleMeasurementsAndUpdateInPlace(measurements); + + assertEquals("path_utils_test_s1", measurements.get(0)); + assertEquals("path_utils_test_s2", measurements.get(1)); + assertEquals(2, measurements.size()); + } + + @Test + public void testCheckSingleMeasurementListsInPlaceReusesStrings() throws MetadataException { + List> measurementLists = new ArrayList<>(); + measurementLists.add( + new ArrayList<>( + Arrays.asList(new String("path_utils_batch_s1"), new String("`path_utils_batch_s2`")))); + measurementLists.add( + new ArrayList<>( + Arrays.asList(new String("path_utils_batch_s1"), new String("`path_utils_batch_s2`")))); + + PathUtils.checkIsLegalSingleMeasurementListsAndUpdateInPlace(measurementLists); + + assertSame(measurementLists.get(0).get(0), measurementLists.get(1).get(0)); + assertSame(measurementLists.get(0).get(1), measurementLists.get(1).get(1)); + assertEquals("path_utils_batch_s2", measurementLists.get(0).get(1)); + } + + @Test + public void testCheckSingleMeasurementListInPlace() throws MetadataException { + List measurements = + new ArrayList<>(Arrays.asList("path_utils_batch_s1", "`path_utils_batch_s2`", null)); + + PathUtils.checkIsLegalSingleMeasurementListsAndUpdateInPlace( + Collections.singletonList(measurements)); + + assertEquals("path_utils_batch_s1", measurements.get(0)); + assertEquals("path_utils_batch_s2", measurements.get(1)); + assertNull(measurements.get(2)); + } +}