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 e9bbc8e7944 Fix table-model Pipe batch isolation (#18577)
e9bbc8e7944 is described below
commit e9bbc8e79447d736f42a7c2c00da69611dcad14e
Author: Caideyipi <[email protected]>
AuthorDate: Tue Sep 8 09:16:06 2026 +0800
Fix table-model Pipe batch isolation (#18577)
---
.../request/PipeTransferTabletBatchReqV2.java | 39 +++---
.../db/pipe/sink/util/cacher/LeaderCacheUtils.java | 6 +-
.../plan/analyze/schema/SchemaValidator.java | 37 +++++-
.../plan/relational/planner/TableModelPlanner.java | 6 +-
.../pipe/sink/PipeDataNodeThriftRequestTest.java | 133 +++++++++++++++++++++
.../sink/util/cacher/LeaderCacheUtilsTest.java | 62 ++++++++++
.../plan/analyze/schema/SchemaValidatorTest.java | 101 ++++++++++++++++
.../relational/planner/TableModelPlannerTest.java | 78 ++++++++++++
8 files changed, 438 insertions(+), 24 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
index c79cab7d88f..2cec219fa75 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
@@ -64,8 +64,9 @@ public class PipeTransferTabletBatchReqV2 extends
TPipeTransferReq {
final List<InsertBaseStatement> statements =
new ArrayList<>(insertNodeReqs.size() + tabletReqs.size());
- final Map<String, List<InsertRowStatement>>
tableModelDatabaseInsertRowStatementMap =
- new LinkedHashMap<>();
+ // Keep permission checks, schema validation, and redirect metadata scoped
to one table.
+ final Map<String, Map<String, List<InsertRowStatement>>>
+ tableModelDatabaseInsertRowStatementMap = new LinkedHashMap<>();
final Map<String, List<InsertRowStatement>>
treeModelDatabaseInsertRowStatementMap =
new LinkedHashMap<>();
final Map<String, List<InsertTabletStatement>>
treeModelDatabaseInsertTabletStatementMap =
@@ -78,17 +79,15 @@ public class PipeTransferTabletBatchReqV2 extends
TPipeTransferReq {
}
if (statement.isWriteToTable()) {
if (statement instanceof InsertRowStatement) {
- tableModelDatabaseInsertRowStatementMap
- .computeIfAbsent(statement.getDatabaseName().get(), k -> new
ArrayList<>())
- .add((InsertRowStatement) statement);
+ addTableModelInsertRowStatement(
+ tableModelDatabaseInsertRowStatementMap, (InsertRowStatement)
statement);
} else if (statement instanceof InsertTabletStatement) {
statements.add(statement);
} else if (statement instanceof InsertRowsStatement) {
for (final InsertRowStatement insertRowStatement :
((InsertRowsStatement) statement).getInsertRowStatementList()) {
- tableModelDatabaseInsertRowStatementMap
- .computeIfAbsent(insertRowStatement.getDatabaseName().get(), k
-> new ArrayList<>())
- .add(insertRowStatement);
+ addTableModelInsertRowStatement(
+ tableModelDatabaseInsertRowStatementMap, insertRowStatement);
}
} else {
throw new UnsupportedOperationException(
@@ -141,18 +140,30 @@ public class PipeTransferTabletBatchReqV2 extends
TPipeTransferReq {
addTreeModelInsertRowsStatements(statements,
treeModelDatabaseInsertRowStatementMap);
addTreeModelInsertTabletsStatements(statements,
treeModelDatabaseInsertTabletStatementMap);
- for (final Map.Entry<String, List<InsertRowStatement>> insertRows :
+ for (final Map.Entry<String, Map<String, List<InsertRowStatement>>>
insertRows :
tableModelDatabaseInsertRowStatementMap.entrySet()) {
- final InsertRowsStatement statement = new InsertRowsStatement();
- statement.setWriteToTable(true);
- statement.setDatabaseName(insertRows.getKey());
- statement.setInsertRowStatementList(insertRows.getValue());
- statements.add(statement);
+ for (final Map.Entry<String, List<InsertRowStatement>> tableInsertRows :
+ insertRows.getValue().entrySet()) {
+ final InsertRowsStatement statement = new InsertRowsStatement();
+ statement.setWriteToTable(true);
+ statement.setDatabaseName(insertRows.getKey());
+ statement.setInsertRowStatementList(tableInsertRows.getValue());
+ statements.add(statement);
+ }
}
return statements;
}
+ private static void addTableModelInsertRowStatement(
+ final Map<String, Map<String, List<InsertRowStatement>>>
databaseInsertRowStatementMap,
+ final InsertRowStatement insertRowStatement) {
+ databaseInsertRowStatementMap
+ .computeIfAbsent(insertRowStatement.getDatabaseName().get(), k -> new
LinkedHashMap<>())
+ .computeIfAbsent(insertRowStatement.getTableName(), k -> new
ArrayList<>())
+ .add(insertRowStatement);
+ }
+
private void addTreeModelInsertRowsStatements(
final List<InsertBaseStatement> statements,
final Map<String, List<InsertRowStatement>>
databaseInsertRowStatementMap) {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
index 754b406dcd7..0f6beade80d 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
@@ -41,11 +41,11 @@ public class LeaderCacheUtils {
* @return a list of pairs, each pair contains a device path and its
redirect endpoint.
*/
public static List<Pair<String, TEndPoint>>
parseRecommendedRedirections(TSStatus status) {
- // If there is no exception, there should be 2 sub-statuses, one for
InsertRowsStatement and one
- // for InsertMultiTabletsStatement (see
IoTDBDataNodeReceiver#handleTransferTabletBatch).
+ // Each top-level sub-status corresponds to one statement constructed by
the receiver. V2 batch
+ // requests may contain any number of statements because rows are grouped
by database and table.
final List<Pair<String, TEndPoint>> redirectList = new ArrayList<>();
- if (status.getSubStatusSize() != 2) {
+ if (!status.isSetSubStatus()) {
return redirectList;
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java
index 6123afb1d9e..06905feead4 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java
@@ -25,11 +25,14 @@ import org.apache.iotdb.commons.path.PartialPath;
import
org.apache.iotdb.commons.queryengine.plan.relational.metadata.QualifiedObjectName;
import org.apache.iotdb.db.queryengine.common.MPPQueryContext;
import org.apache.iotdb.db.queryengine.common.schematree.ISchemaTree;
+import org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeUtils;
import org.apache.iotdb.db.queryengine.plan.relational.metadata.Metadata;
import org.apache.iotdb.db.queryengine.plan.relational.security.AccessControl;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.InsertRows;
import
org.apache.iotdb.db.queryengine.plan.relational.sql.ast.WrappedInsertStatement;
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
import
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertMultiTabletsStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement;
import
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsOfOneDeviceStatement;
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
@@ -39,9 +42,12 @@ import org.apache.tsfile.file.metadata.enums.TSEncoding;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.LinkedHashSet;
import java.util.List;
+import java.util.Set;
import static org.apache.iotdb.commons.utils.PathUtils.unQualifyDatabaseName;
+import static
org.apache.iotdb.db.queryengine.plan.execution.config.TableConfigTaskVisitor.DATABASE_NOT_SPECIFIED;
public class SchemaValidator {
@@ -71,11 +77,10 @@ public class SchemaValidator {
final MPPQueryContext context,
AccessControl accessControl) {
try {
- accessControl.checkCanInsertIntoTable(
- context.getSession().getUserName(),
- new QualifiedObjectName(
- unQualifyDatabaseName(insertStatement.getDatabase()),
insertStatement.getTableName()),
- context);
+ for (final QualifiedObjectName targetTable :
getTargetTables(insertStatement, context)) {
+ accessControl.checkCanInsertIntoTable(
+ context.getSession().getUserName(), targetTable, context);
+ }
insertStatement.validateTableSchema(metadata, context);
insertStatement.updateAfterSchemaValidation(context);
insertStatement.validateDeviceSchema(metadata, context);
@@ -85,6 +90,28 @@ public class SchemaValidator {
}
}
+ private static Set<QualifiedObjectName> getTargetTables(
+ final WrappedInsertStatement insertStatement, final MPPQueryContext
context) {
+ final Set<QualifiedObjectName> targetTables = new LinkedHashSet<>();
+ if (insertStatement instanceof InsertRows) {
+ for (final InsertRowStatement rowStatement :
+ ((InsertRows)
insertStatement).getInnerTreeStatement().getInsertRowStatementList()) {
+ final String database = AnalyzeUtils.getDatabaseName(rowStatement,
context);
+ if (database == null) {
+ throw new SemanticException(DATABASE_NOT_SPECIFIED);
+ }
+ targetTables.add(
+ new QualifiedObjectName(unQualifyDatabaseName(database),
rowStatement.getTableName()));
+ }
+ } else {
+ targetTables.add(
+ new QualifiedObjectName(
+ unQualifyDatabaseName(insertStatement.getDatabase()),
+ insertStatement.getTableName()));
+ }
+ return targetTables;
+ }
+
public static ISchemaTree validate(
ISchemaFetcher schemaFetcher,
List<PartialPath> devicePaths,
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlanner.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlanner.java
index c8433d95416..cca616e0a3b 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlanner.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlanner.java
@@ -55,6 +55,7 @@ import
org.apache.iotdb.db.queryengine.plan.scheduler.ClusterScheduler;
import org.apache.iotdb.db.queryengine.plan.scheduler.IScheduler;
import org.apache.iotdb.db.queryengine.plan.scheduler.load.LoadTsFileScheduler;
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
import
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertTabletStatement;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -238,8 +239,9 @@ public class TableModelPlanner implements IPlanner {
((WrappedInsertStatement) statementToRedirect).getInnerTreeStatement();
if (!analysis.isFinishQueryAfterAnalyze()) {
- // Table Model Session only supports insertTablet
- if (insertStatement instanceof InsertTabletStatement) {
+ // Table Model Session supports insertTablet and pipe-generated
insertRows statements.
+ if (insertStatement instanceof InsertTabletStatement
+ || insertStatement instanceof InsertRowsStatement) {
if (tsstatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode())
{
boolean needRedirect = false;
List<TEndPoint> redirectNodeList = analysis.getRedirectNodeList();
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
index 143d195d284..b780d19f01c 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
@@ -53,6 +53,7 @@ import
org.apache.iotdb.db.pipe.sink.payload.evolvable.request.PipeTransferTsFil
import
org.apache.iotdb.db.pipe.sink.payload.evolvable.request.PipeTransferTsFileSealWithModReq;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.CreateAlignedTimeSeriesNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
+import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode;
import org.apache.iotdb.db.queryengine.plan.statement.Statement;
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
import
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertMultiTabletsStatement;
@@ -1048,6 +1049,138 @@ public class PipeDataNodeThriftRequestTest {
new HashSet<>(java.util.Arrays.asList("root.db1", "root.db2")),
insertTabletsDatabases);
}
+ @Test
+ public void testPipeTransferTabletBatchReqV2SeparatesTableModelTables()
throws IOException {
+ final List<ByteBuffer> insertNodeBuffers = new ArrayList<>();
+ final List<String> insertNodeDataBases = new ArrayList<>();
+
+ insertNodeBuffers.add(
+ new InsertRowNode(
+ new PlanNodeId(""),
+ new PartialPath("table1", false),
+ false,
+ new String[] {"s"},
+ new TSDataType[] {TSDataType.INT32},
+ 1,
+ new Object[] {1},
+ false)
+ .serializeToByteBuffer());
+ insertNodeDataBases.add("db1");
+
+ insertNodeBuffers.add(
+ new InsertRowNode(
+ new PlanNodeId(""),
+ new PartialPath("table2", false),
+ false,
+ new String[] {"s"},
+ new TSDataType[] {TSDataType.INT32},
+ 2,
+ new Object[] {2},
+ false)
+ .serializeToByteBuffer());
+ insertNodeDataBases.add("db1");
+
+ insertNodeBuffers.add(
+ new InsertRowNode(
+ new PlanNodeId(""),
+ new PartialPath("table1", false),
+ false,
+ new String[] {"s"},
+ new TSDataType[] {TSDataType.INT32},
+ 3,
+ new Object[] {3},
+ false)
+ .serializeToByteBuffer());
+ insertNodeDataBases.add("db1");
+
+ final PipeTransferTabletBatchReqV2 request =
+ PipeTransferTabletBatchReqV2.fromTPipeTransferReq(
+ PipeTransferTabletBatchReqV2.toTPipeTransferReq(
+ insertNodeBuffers,
+ Collections.emptyList(),
+ insertNodeDataBases,
+ Collections.emptyList()));
+
+ final List<InsertBaseStatement> statements = request.constructStatements();
+
+ Assert.assertEquals(2, statements.size());
+ final InsertRowsStatement table1Statement = (InsertRowsStatement)
statements.get(0);
+ final InsertRowsStatement table2Statement = (InsertRowsStatement)
statements.get(1);
+ Assert.assertTrue(table1Statement.isWriteToTable());
+ Assert.assertTrue(table2Statement.isWriteToTable());
+ Assert.assertEquals("db1", table1Statement.getDatabaseName().get());
+ Assert.assertEquals("db1", table2Statement.getDatabaseName().get());
+ Assert.assertEquals(2, table1Statement.getInsertRowStatementList().size());
+ Assert.assertEquals(1, table2Statement.getInsertRowStatementList().size());
+ Assert.assertEquals(
+ "table1",
table1Statement.getInsertRowStatementList().get(0).getTableName());
+ Assert.assertEquals(
+ "table1",
table1Statement.getInsertRowStatementList().get(1).getTableName());
+ Assert.assertEquals(
+ "table2",
table2Statement.getInsertRowStatementList().get(0).getTableName());
+ }
+
+ @Test
+ public void
testPipeTransferTabletBatchReqV2SeparatesTablesWithinInsertRowsNode()
+ throws IOException {
+ final InsertRowsNode insertRowsNode = new InsertRowsNode(new
PlanNodeId("rows"));
+ insertRowsNode.addOneInsertRowNode(
+ new InsertRowNode(
+ new PlanNodeId("row1"),
+ new PartialPath("table1", false),
+ false,
+ new String[] {"s"},
+ new TSDataType[] {TSDataType.INT32},
+ 1,
+ new Object[] {1},
+ false),
+ 0);
+ insertRowsNode.addOneInsertRowNode(
+ new InsertRowNode(
+ new PlanNodeId("row2"),
+ new PartialPath("table2", false),
+ false,
+ new String[] {"s"},
+ new TSDataType[] {TSDataType.INT32},
+ 2,
+ new Object[] {2},
+ false),
+ 1);
+ insertRowsNode.addOneInsertRowNode(
+ new InsertRowNode(
+ new PlanNodeId("row3"),
+ new PartialPath("table1", false),
+ false,
+ new String[] {"s"},
+ new TSDataType[] {TSDataType.INT32},
+ 3,
+ new Object[] {3},
+ false),
+ 2);
+
+ final PipeTransferTabletBatchReqV2 request =
+ PipeTransferTabletBatchReqV2.fromTPipeTransferReq(
+ PipeTransferTabletBatchReqV2.toTPipeTransferReq(
+
Collections.singletonList(insertRowsNode.serializeToByteBuffer()),
+ Collections.emptyList(),
+ Collections.singletonList("db1"),
+ Collections.emptyList()));
+
+ final List<InsertBaseStatement> statements = request.constructStatements();
+
+ Assert.assertEquals(2, statements.size());
+ final InsertRowsStatement table1Statement = (InsertRowsStatement)
statements.get(0);
+ final InsertRowsStatement table2Statement = (InsertRowsStatement)
statements.get(1);
+ Assert.assertEquals(2, table1Statement.getInsertRowStatementList().size());
+ Assert.assertEquals(1, table2Statement.getInsertRowStatementList().size());
+ Assert.assertEquals(
+ "table1",
table1Statement.getInsertRowStatementList().get(0).getTableName());
+ Assert.assertEquals(
+ "table1",
table1Statement.getInsertRowStatementList().get(1).getTableName());
+ Assert.assertEquals(
+ "table2",
table2Statement.getInsertRowStatementList().get(0).getTableName());
+ }
+
@Test
public void testPipeTransferFilePieceReq() throws IOException {
final byte[] body = "testPipeTransferFilePieceReq".getBytes();
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
new file mode 100644
index 00000000000..76c6bef2467
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
@@ -0,0 +1,62 @@
+/*
+ * 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.util.cacher;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.rpc.RpcUtils;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.apache.tsfile.utils.Pair;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+public class LeaderCacheUtilsTest {
+
+ @Test
+ public void testParseRecommendedRedirectionsFromVariableStatementCount() {
+ final TEndPoint redirectEndPoint = new TEndPoint("127.0.0.2", 6667);
+ final TSStatus redirectedRowStatus =
+ RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS)
+ .setMessage("table1.device1")
+ .setRedirectNode(redirectEndPoint);
+ final TSStatus redirectedStatementStatus =
+ RpcUtils.getStatus(TSStatusCode.REDIRECTION_RECOMMEND)
+ .setSubStatus(Collections.singletonList(redirectedRowStatus));
+ final TSStatus batchStatus =
+ RpcUtils.getStatus(TSStatusCode.REDIRECTION_RECOMMEND)
+ .setSubStatus(
+ Arrays.asList(
+ RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS),
+ redirectedStatementStatus,
+ RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS)));
+
+ final List<Pair<String, TEndPoint>> redirects =
+ LeaderCacheUtils.parseRecommendedRedirections(batchStatus);
+
+ Assert.assertEquals(1, redirects.size());
+ Assert.assertEquals("table1.device1", redirects.get(0).getLeft());
+ Assert.assertEquals(redirectEndPoint, redirects.get(0).getRight());
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidatorTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidatorTest.java
new file mode 100644
index 00000000000..8cb3cfb1ee1
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidatorTest.java
@@ -0,0 +1,101 @@
+/*
+ * 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.queryengine.plan.analyze.schema;
+
+import org.apache.iotdb.commons.exception.auth.AccessDeniedException;
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.queryengine.common.SessionInfo;
+import org.apache.iotdb.commons.queryengine.common.SqlDialect;
+import
org.apache.iotdb.commons.queryengine.plan.relational.metadata.QualifiedObjectName;
+import org.apache.iotdb.db.queryengine.common.MPPQueryContext;
+import org.apache.iotdb.db.queryengine.common.QueryId;
+import org.apache.iotdb.db.queryengine.plan.relational.metadata.Metadata;
+import org.apache.iotdb.db.queryengine.plan.relational.security.AccessControl;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.InsertRows;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
+
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.InOrder;
+import org.mockito.Mockito;
+
+import java.time.ZoneId;
+import java.util.Arrays;
+
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.same;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
+import static org.mockito.Mockito.verifyZeroInteractions;
+
+public class SchemaValidatorTest {
+
+ private static final String USER = "pipe-user";
+
+ @Test
+ public void testAllInsertRowsTablesCheckedBeforeSchemaValidation() {
+ final MPPQueryContext context =
+ new MPPQueryContext(
+ "",
+ new QueryId("check_all_insert_rows_tables"),
+ new SessionInfo(1L, USER, ZoneId.systemDefault(), "db1",
SqlDialect.TABLE),
+ null,
+ null);
+ final InsertRowsStatement statement = new InsertRowsStatement();
+ statement.setWriteToTable(true);
+ statement.setInsertRowStatementList(
+ Arrays.asList(
+ createInsertRowStatement("db1", "table1"),
+ createInsertRowStatement("db1", "table2"),
+ createInsertRowStatement("db1", "table1")));
+
+ final Metadata metadata = Mockito.mock(Metadata.class);
+ final AccessControl accessControl = Mockito.mock(AccessControl.class);
+ final QualifiedObjectName table1 = new QualifiedObjectName("db1",
"table1");
+ final QualifiedObjectName table2 = new QualifiedObjectName("db1",
"table2");
+ Mockito.doThrow(new AccessDeniedException("denied"))
+ .when(accessControl)
+ .checkCanInsertIntoTable(eq(USER), eq(table2), same(context));
+
+ Assert.assertThrows(
+ AccessDeniedException.class,
+ () ->
+ SchemaValidator.validate(
+ metadata, new InsertRows(statement, context), context,
accessControl));
+
+ final InOrder inOrder = Mockito.inOrder(accessControl);
+ inOrder.verify(accessControl).checkCanInsertIntoTable(eq(USER),
eq(table1), same(context));
+ inOrder.verify(accessControl).checkCanInsertIntoTable(eq(USER),
eq(table2), same(context));
+ verify(accessControl, times(1)).checkCanInsertIntoTable(eq(USER),
eq(table1), same(context));
+ verifyNoMoreInteractions(accessControl);
+ verifyZeroInteractions(metadata);
+ }
+
+ private static InsertRowStatement createInsertRowStatement(
+ final String database, final String table) {
+ final InsertRowStatement statement = new InsertRowStatement();
+ statement.setWriteToTable(true);
+ statement.setDatabaseName(database);
+ statement.setDevicePath(new PartialPath(table, false));
+ return statement;
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlannerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlannerTest.java
new file mode 100644
index 00000000000..99c0ed02357
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlannerTest.java
@@ -0,0 +1,78 @@
+/*
+ * 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.queryengine.plan.relational.planner;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.db.queryengine.plan.relational.analyzer.Analysis;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.InsertRows;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
+import org.apache.iotdb.rpc.RpcUtils;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+public class TableModelPlannerTest {
+
+ @Test
+ public void testSetRedirectInfoForInsertRows() {
+ final InsertRowsStatement insertRowsStatement = new InsertRowsStatement();
+ insertRowsStatement.setInsertRowStatementList(Collections.emptyList());
+ final Analysis analysis =
+ new Analysis(new InsertRows(insertRowsStatement, null),
Collections.emptyMap());
+ final TEndPoint localEndPoint = new TEndPoint("127.0.0.1", 6667);
+ final TEndPoint remoteEndPoint = new TEndPoint("127.0.0.2", 6667);
+ analysis.setRedirectNodeList(Arrays.asList(localEndPoint, remoteEndPoint));
+ final TSStatus status = RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS);
+
+ createPlanner().setRedirectInfo(analysis, localEndPoint, status);
+
+ Assert.assertEquals(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode(),
status.getCode());
+ final List<TSStatus> subStatus = status.getSubStatus();
+ Assert.assertEquals(2, subStatus.size());
+ Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(),
subStatus.get(0).getCode());
+ Assert.assertFalse(subStatus.get(0).isSetRedirectNode());
+ Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(),
subStatus.get(1).getCode());
+ Assert.assertEquals(remoteEndPoint, subStatus.get(1).getRedirectNode());
+ }
+
+ private static TableModelPlanner createPlanner() {
+ return new TableModelPlanner(
+ null,
+ null,
+ null,
+ null,
+ null,
+ null,
+ null,
+ Collections.emptyList(),
+ Collections.emptyList(),
+ null,
+ null,
+ Collections.emptyList(),
+ Collections.emptyMap(),
+ null);
+ }
+}