Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -64,8 +64,9 @@ public List<InsertBaseStatement> constructStatements() {
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 =
Expand All @@ -78,17 +79,15 @@ public List<InsertBaseStatement> constructStatements() {
}
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(
Expand Down Expand Up @@ -141,18 +140,30 @@ public List<InsertBaseStatement> constructStatements() {
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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,11 +41,11 @@ private 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;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,11 +25,14 @@
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;

Expand All @@ -39,9 +42,12 @@
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 {

Expand Down Expand Up @@ -71,11 +77,10 @@ public static void validate(
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);
Expand All @@ -85,6 +90,28 @@ public static void validate(
}
}

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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@
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;
Expand Down Expand Up @@ -238,8 +239,9 @@ public void setRedirectInfo(IAnalysis iAnalysis, TEndPoint localEndPoint, TSStat
((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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
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;
Expand Down Expand Up @@ -1048,6 +1049,138 @@ public void testPipeTransferTabletBatchReqV2WithMultipleTreeModelDatabases() thr
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();
Expand Down
Original file line number Diff line number Diff line change
@@ -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());
}
}
Loading
Loading