Skip to content

Commit ef45dbf

Browse files
CaideyipiJackieTien97
authored andcommitted
Fix pipe tree database creation on receiver (#17991)
1 parent 9a310b6 commit ef45dbf

5 files changed

Lines changed: 239 additions & 8 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java‎

Lines changed: 144 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,8 @@
2626
import org.apache.iotdb.commons.disk.strategy.DirectoryStrategyType;
2727
import org.apache.iotdb.commons.exception.DiskSpaceInsufficientException;
2828
import org.apache.iotdb.commons.exception.IllegalPathException;
29+
import org.apache.iotdb.commons.exception.IoTDBException;
30+
import org.apache.iotdb.commons.exception.IoTDBRuntimeException;
2931
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeOutOfMemoryCriticalException;
3032
import org.apache.iotdb.commons.path.PartialPath;
3133
import org.apache.iotdb.commons.pipe.config.PipeConfig;
@@ -89,6 +91,7 @@
8991
import org.apache.iotdb.db.queryengine.plan.analyze.schema.ClusterSchemaFetcher;
9092
import org.apache.iotdb.db.queryengine.plan.execution.config.ConfigTaskResult;
9193
import org.apache.iotdb.db.queryengine.plan.execution.config.executor.ClusterConfigTaskExecutor;
94+
import org.apache.iotdb.db.queryengine.plan.execution.config.metadata.DatabaseSchemaTask;
9295
import org.apache.iotdb.db.queryengine.plan.execution.config.metadata.relational.CreateDBTask;
9396
import org.apache.iotdb.db.queryengine.plan.planner.LocalExecutionPlanner;
9497
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.view.AlterLogicalViewNode;
@@ -99,6 +102,7 @@
99102
import org.apache.iotdb.db.queryengine.plan.statement.StatementType;
100103
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
101104
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertMultiTabletsStatement;
105+
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsOfOneDeviceStatement;
102106
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
103107
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertTabletStatement;
104108
import org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement;
@@ -133,6 +137,7 @@
133137
import java.util.Objects;
134138
import java.util.Optional;
135139
import java.util.Set;
140+
import java.util.concurrent.ConcurrentHashMap;
136141
import java.util.concurrent.ExecutionException;
137142
import java.util.concurrent.atomic.AtomicLong;
138143
import java.util.concurrent.atomic.AtomicReference;
@@ -166,7 +171,8 @@ public class IoTDBDataNodeReceiver extends IoTDBFileReceiver {
166171
this::executeStatementForTableModel);
167172
private final PipeTreeStatementDataTypeConvertExecutionVisitor
168173
treeStatementDataTypeConvertExecutionVisitor =
169-
new PipeTreeStatementDataTypeConvertExecutionVisitor(this::executeStatementForTreeModel);
174+
new PipeTreeStatementDataTypeConvertExecutionVisitor(
175+
statement -> executeStatementForTreeModel(statement, getTreeDatabaseName(statement)));
170176
public final PipeTreeStatementToBatchVisitor batchVisitor = new PipeTreeStatementToBatchVisitor();
171177

172178
// Used for data transfer: confignode (cluster A) -> datanode (cluster B) -> confignode (cluster
@@ -186,6 +192,14 @@ public class IoTDBDataNodeReceiver extends IoTDBFileReceiver {
186192
private static final PipeConfig PIPE_CONFIG = PipeConfig.getInstance();
187193

188194
private PipeMemoryBlock allocatedMemoryBlock;
195+
private final Set<String> autoCreatedTreeDatabases = ConcurrentHashMap.newKeySet();
196+
private final Set<String> conflictedTreeDatabases = ConcurrentHashMap.newKeySet();
197+
198+
private enum TreeDatabaseCreationResult {
199+
SKIPPED,
200+
CREATED_OR_EXISTED,
201+
CONFLICTED
202+
}
189203

190204
static {
191205
try {
@@ -965,6 +979,9 @@ private TSStatus executeStatementWithPermissionCheckAndRetryOnDataTypeMismatch(
965979
((InsertBaseStatement) statement).getDatabaseName().isPresent()
966980
? ((InsertBaseStatement) statement).getDatabaseName().get()
967981
: null;
982+
} else if (statement instanceof InsertBaseStatement) {
983+
isTableModelStatement = false;
984+
databaseName = getTreeDatabaseName(statement);
968985
} else {
969986
isTableModelStatement = false;
970987
databaseName = null;
@@ -1015,7 +1032,7 @@ private TSStatus executeStatementWithPermissionCheckAndRetryOnDataTypeMismatch(
10151032
final TSStatus status =
10161033
isTableModelStatement
10171034
? executeStatementForTableModel(statement, databaseName)
1018-
: executeStatementForTreeModel(statement);
1035+
: executeStatementForTreeModel(statement, getTreeDatabaseName(statement));
10191036

10201037
// Try to convert data type if the status code is not success. Insert statements normally return
10211038
// above after the first converted execution. The retry path is kept for load and fallback
@@ -1158,7 +1175,84 @@ private void autoCreateDatabaseIfNecessary(final String database) {
11581175
}
11591176
}
11601177

1161-
private TSStatus executeStatementForTreeModel(final Statement statement) {
1178+
private TreeDatabaseCreationResult autoCreateTreeDatabaseIfNecessary(final String database) {
1179+
if (database == null
1180+
|| LoadTsFileStatement.getDatabaseLevelByTreeDatabase(database) == null
1181+
|| !IoTDBDescriptor.getInstance().getConfig().isAutoCreateSchemaEnabled()) {
1182+
return TreeDatabaseCreationResult.SKIPPED;
1183+
}
1184+
if (autoCreatedTreeDatabases.contains(database)) {
1185+
return TreeDatabaseCreationResult.CREATED_OR_EXISTED;
1186+
}
1187+
if (conflictedTreeDatabases.contains(database)) {
1188+
return TreeDatabaseCreationResult.CONFLICTED;
1189+
}
1190+
1191+
try {
1192+
final TSStatus status =
1193+
AuthorityChecker.getAccessControl()
1194+
.checkCanCreateDatabaseForTree(getUserEntity(), new PartialPath(database));
1195+
if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
1196+
throw new PipeException(status.getMessage());
1197+
}
1198+
1199+
final DatabaseSchemaStatement statement =
1200+
new DatabaseSchemaStatement(DatabaseSchemaStatement.DatabaseSchemaStatementType.CREATE);
1201+
statement.setDatabasePath(new PartialPath(database));
1202+
statement.setEnablePrintExceptionLog(false);
1203+
final DatabaseSchemaTask task = new DatabaseSchemaTask(statement);
1204+
final ListenableFuture<ConfigTaskResult> future =
1205+
task.execute(ClusterConfigTaskExecutor.getInstance());
1206+
final ConfigTaskResult result = future.get();
1207+
final int statusCode = result.getStatusCode().getStatusCode();
1208+
if (statusCode == TSStatusCode.SUCCESS_STATUS.getStatusCode()
1209+
|| statusCode == TSStatusCode.DATABASE_ALREADY_EXISTS.getStatusCode()) {
1210+
autoCreatedTreeDatabases.add(database);
1211+
return TreeDatabaseCreationResult.CREATED_OR_EXISTED;
1212+
}
1213+
if (statusCode == TSStatusCode.DATABASE_CONFLICT.getStatusCode()) {
1214+
conflictedTreeDatabases.add(database);
1215+
return TreeDatabaseCreationResult.CONFLICTED;
1216+
}
1217+
throw new PipeException(
1218+
String.format(
1219+
"Auto create tree database failed: %s, status code: %s",
1220+
database, result.getStatusCode()));
1221+
} catch (final IllegalPathException e) {
1222+
throw new PipeException(String.format("Illegal tree database %s.", database), e);
1223+
} catch (final ExecutionException | InterruptedException e) {
1224+
if (e instanceof InterruptedException) {
1225+
Thread.currentThread().interrupt();
1226+
}
1227+
final Throwable rootCause = getRootCause(e);
1228+
final int errorCode;
1229+
if (rootCause instanceof IoTDBException) {
1230+
errorCode = ((IoTDBException) rootCause).getErrorCode();
1231+
} else if (rootCause instanceof IoTDBRuntimeException) {
1232+
errorCode = ((IoTDBRuntimeException) rootCause).getErrorCode();
1233+
} else {
1234+
errorCode = TSStatusCode.INTERNAL_SERVER_ERROR.getStatusCode();
1235+
}
1236+
if (errorCode == TSStatusCode.DATABASE_ALREADY_EXISTS.getStatusCode()) {
1237+
autoCreatedTreeDatabases.add(database);
1238+
return TreeDatabaseCreationResult.CREATED_OR_EXISTED;
1239+
}
1240+
if (errorCode == TSStatusCode.DATABASE_CONFLICT.getStatusCode()) {
1241+
conflictedTreeDatabases.add(database);
1242+
return TreeDatabaseCreationResult.CONFLICTED;
1243+
}
1244+
throw new PipeException(
1245+
DataNodePipeMessages.AUTO_CREATE_DATABASE_FAILED_BECAUSE + e.getMessage());
1246+
}
1247+
}
1248+
1249+
private TSStatus executeStatementForTreeModel(
1250+
final Statement statement, final String databaseName) {
1251+
if (autoCreateTreeDatabaseIfNecessary(databaseName) == TreeDatabaseCreationResult.CONFLICTED) {
1252+
// Continue execution, but let partition analysis infer the receiver-side database.
1253+
clearTreeDatabaseName(statement);
1254+
}
1255+
11621256
return Coordinator.getInstance()
11631257
.executeForTreeModel(
11641258
shouldMarkAsPipeRequest.get() ? new PipeEnrichedStatement(statement) : statement,
@@ -1173,6 +1267,53 @@ private TSStatus executeStatementForTreeModel(final Statement statement) {
11731267
.status;
11741268
}
11751269

1270+
private IAuditEntity getUserEntity() {
1271+
return userEntity != null
1272+
? userEntity
1273+
: AuthorityChecker.createIAuditEntity(username, SESSION_MANAGER.getCurrSession());
1274+
}
1275+
1276+
private String getTreeDatabaseName(final Statement statement) {
1277+
if (statement instanceof LoadTsFileStatement) {
1278+
return ((LoadTsFileStatement) statement).getDatabase();
1279+
}
1280+
if (statement instanceof InsertBaseStatement) {
1281+
return ((InsertBaseStatement) statement).getDatabaseName().orElse(null);
1282+
}
1283+
return null;
1284+
}
1285+
1286+
static void clearTreeDatabaseName(final Statement statement) {
1287+
if (statement instanceof LoadTsFileStatement) {
1288+
final LoadTsFileStatement loadTsFileStatement = (LoadTsFileStatement) statement;
1289+
loadTsFileStatement.setDatabase(null);
1290+
loadTsFileStatement.setDatabaseLevel(
1291+
IoTDBDescriptor.getInstance().getConfig().getDefaultDatabaseLevel());
1292+
} else if (statement instanceof InsertBaseStatement) {
1293+
clearTreeInsertDatabaseName((InsertBaseStatement) statement);
1294+
}
1295+
}
1296+
1297+
private static void clearTreeInsertDatabaseName(final InsertBaseStatement statement) {
1298+
statement.setDatabaseName(null);
1299+
if (statement instanceof InsertRowsStatement) {
1300+
for (final InsertBaseStatement childStatement :
1301+
((InsertRowsStatement) statement).getInsertRowStatementList()) {
1302+
childStatement.setDatabaseName(null);
1303+
}
1304+
} else if (statement instanceof InsertRowsOfOneDeviceStatement) {
1305+
for (final InsertBaseStatement childStatement :
1306+
((InsertRowsOfOneDeviceStatement) statement).getInsertRowStatementList()) {
1307+
childStatement.setDatabaseName(null);
1308+
}
1309+
} else if (statement instanceof InsertMultiTabletsStatement) {
1310+
for (final InsertBaseStatement childStatement :
1311+
((InsertMultiTabletsStatement) statement).getInsertTabletStatementList()) {
1312+
childStatement.setDatabaseName(null);
1313+
}
1314+
}
1315+
}
1316+
11761317
private TSStatus executeStatementForTableModelWithPermissionCheck(
11771318
final org.apache.iotdb.commons.queryengine.plan.relational.sql.ast.Statement statement,
11781319
final String databaseName) {

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/visitor/PipeTreeStatementDataTypeConvertExecutionVisitor.java‎

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -105,12 +105,15 @@ public Optional<TSStatus> visitLoadFile(
105105
new TsFileInsertionEventScanParser(
106106
file, new IoTDBTreePattern(null), Long.MIN_VALUE, Long.MAX_VALUE, null, null, true)) {
107107
for (final Pair<Tablet, Boolean> tabletWithIsAligned : parser.toTabletWithIsAligneds()) {
108+
final InsertTabletStatement insertTabletStatement =
109+
PipeTransferTabletRawReq.toTPipeTransferRawReq(
110+
tabletWithIsAligned.getLeft(), tabletWithIsAligned.getRight())
111+
.constructStatement();
112+
if (loadTsFileStatement.getDatabase() != null) {
113+
insertTabletStatement.setDatabaseName(loadTsFileStatement.getDatabase());
114+
}
108115
final PipeConvertedInsertTabletStatement statement =
109-
new PipeConvertedInsertTabletStatement(
110-
PipeTransferTabletRawReq.toTPipeTransferRawReq(
111-
tabletWithIsAligned.getLeft(), tabletWithIsAligned.getRight())
112-
.constructStatement(),
113-
false);
116+
new PipeConvertedInsertTabletStatement(insertTabletStatement, false);
114117

115118
TSStatus result;
116119
try {

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -652,6 +652,7 @@ private LoadTsFileStatement buildRetryTreeLoadStatement(
652652
.setConvertOnTypeMismatch(true);
653653
if (database != null) {
654654
statement.setDatabase(database);
655+
statement.updateDatabaseLevelByTreeDatabase();
655656
}
656657
if (isGeneratedByPipe) {
657658
statement.markIsGeneratedByPipe();

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiverTest.java‎

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,10 @@
2020
package org.apache.iotdb.db.pipe.receiver.protocol.thrift;
2121

2222
import org.apache.iotdb.db.conf.IoTDBDescriptor;
23+
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertMultiTabletsStatement;
24+
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement;
25+
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
26+
import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertTabletStatement;
2327
import org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement;
2428
import org.apache.iotdb.db.storageengine.load.active.ActiveLoadPathHelper;
2529
import org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator;
@@ -29,6 +33,8 @@
2933

3034
import java.nio.file.Files;
3135
import java.nio.file.Path;
36+
import java.util.Arrays;
37+
import java.util.Collections;
3238
import java.util.Map;
3339

3440
public class IoTDBDataNodeReceiverTest {
@@ -140,4 +146,53 @@ public void testLoadTsFileSyncStatementCanSkipVerifySchemaWhenNotConvertingType(
140146
Files.deleteIfExists(tsFile);
141147
}
142148
}
149+
150+
@Test
151+
public void testClearTreeDatabaseNameForLoadTsFileStatement() throws Exception {
152+
final Path tsFile = Files.createTempFile("pipe-load-clear-tree-database", ".tsfile");
153+
try {
154+
final LoadTsFileStatement statement =
155+
IoTDBDataNodeReceiver.buildLoadTsFileStatementForSync(
156+
"root.test.sg_0", tsFile.toString(), true, true);
157+
158+
IoTDBDataNodeReceiver.clearTreeDatabaseName(statement);
159+
160+
Assert.assertNull(statement.getDatabase());
161+
Assert.assertEquals(
162+
IoTDBDescriptor.getInstance().getConfig().getDefaultDatabaseLevel(),
163+
statement.getDatabaseLevel());
164+
} finally {
165+
Files.deleteIfExists(tsFile);
166+
}
167+
}
168+
169+
@Test
170+
public void testClearTreeDatabaseNameForBatchInsertStatements() {
171+
final InsertRowStatement rowStatement1 = new InsertRowStatement();
172+
rowStatement1.setDatabaseName("root.test.sg_0");
173+
final InsertRowStatement rowStatement2 = new InsertRowStatement();
174+
rowStatement2.setDatabaseName("root.test.sg_0");
175+
final InsertRowsStatement insertRowsStatement = new InsertRowsStatement();
176+
insertRowsStatement.setDatabaseName("root.test.sg_0");
177+
insertRowsStatement.setInsertRowStatementList(Arrays.asList(rowStatement1, rowStatement2));
178+
179+
IoTDBDataNodeReceiver.clearTreeDatabaseName(insertRowsStatement);
180+
181+
Assert.assertFalse(insertRowsStatement.getDatabaseName().isPresent());
182+
Assert.assertFalse(rowStatement1.getDatabaseName().isPresent());
183+
Assert.assertFalse(rowStatement2.getDatabaseName().isPresent());
184+
185+
final InsertTabletStatement tabletStatement = new InsertTabletStatement();
186+
tabletStatement.setDatabaseName("root.test.sg_0");
187+
final InsertMultiTabletsStatement insertMultiTabletsStatement =
188+
new InsertMultiTabletsStatement();
189+
insertMultiTabletsStatement.setDatabaseName("root.test.sg_0");
190+
insertMultiTabletsStatement.setInsertTabletStatementList(
191+
Collections.singletonList(tabletStatement));
192+
193+
IoTDBDataNodeReceiver.clearTreeDatabaseName(insertMultiTabletsStatement);
194+
195+
Assert.assertFalse(insertMultiTabletsStatement.getDatabaseName().isPresent());
196+
Assert.assertFalse(tabletStatement.getDatabaseName().isPresent());
197+
}
143198
}

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java‎

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,13 +28,17 @@
2828
import org.apache.iotdb.db.queryengine.plan.planner.plan.PlanFragment;
2929
import org.apache.iotdb.db.queryengine.plan.planner.plan.SubPlan;
3030
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.load.LoadSingleTsFileNode;
31+
import org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement;
3132

3233
import org.junit.Assert;
3334
import org.junit.Before;
3435
import org.junit.Test;
3536
import org.mockito.Mock;
3637
import org.mockito.MockitoAnnotations;
3738

39+
import java.io.File;
40+
import java.lang.reflect.Method;
41+
3842
import static org.mockito.Mockito.mock;
3943
import static org.mockito.Mockito.spy;
4044
import static org.mockito.Mockito.when;
@@ -87,4 +91,31 @@ public void testGetPartitionQueryDatabaseForTableModelLoad() {
8791

8892
Assert.assertEquals("test", LoadTsFileScheduler.getPartitionQueryDatabase(node, false));
8993
}
94+
95+
@Test
96+
public void testBuildRetryTreeLoadStatementUpdatesDatabaseLevel() throws Exception {
97+
final LoadTsFileScheduler scheduler =
98+
new LoadTsFileScheduler(
99+
distributedQueryPlan,
100+
mock(MPPQueryContext.class),
101+
mock(QueryStateMachine.class),
102+
mock(IClientManager.class),
103+
mock(IPartitionFetcher.class),
104+
true);
105+
final Method method =
106+
LoadTsFileScheduler.class.getDeclaredMethod(
107+
"buildRetryTreeLoadStatement", String.class, boolean.class, String.class);
108+
method.setAccessible(true);
109+
110+
final File tsFile = File.createTempFile("test", ".tsfile");
111+
tsFile.deleteOnExit();
112+
113+
final LoadTsFileStatement statement =
114+
(LoadTsFileStatement)
115+
method.invoke(scheduler, tsFile.getAbsolutePath(), true, "root.test.sg_0");
116+
117+
Assert.assertEquals("root.test.sg_0", statement.getDatabase());
118+
Assert.assertEquals(2, statement.getDatabaseLevel());
119+
Assert.assertTrue(statement.isGeneratedByPipe());
120+
}
90121
}

0 commit comments

Comments
 (0)