Skip to content

Commit abcc51c

Browse files
authored
Fix pipe receiver type conversion load path (#17849) (#17873)
Backport receiver-side tree-model changes to dev/1.3. The 2.x table-model load analyzer changes and table/dual ITs are not applicable to dev/1.3. (cherry picked from commit 90055d5)
1 parent 6348c2d commit abcc51c

2 files changed

Lines changed: 65 additions & 11 deletions

File tree

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

Lines changed: 27 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -501,24 +501,31 @@ static Map<String, String> buildLoadTsFileAttributesForAsync(
501501
dataBaseName,
502502
LoadTsFileStatement.getDatabaseLevelByTreeDatabase(dataBaseName),
503503
shouldConvertDataTypeOnTypeMismatch,
504-
validateTsFile,
504+
validateTsFile || shouldConvertDataTypeOnTypeMismatch,
505505
null,
506506
shouldMarkAsPipeRequest);
507507
}
508508

509509
private TSStatus loadTsFileSync(final String dataBaseName, final String fileAbsolutePath)
510510
throws FileNotFoundException {
511511
return executeStatementAndClassifyExceptions(
512-
buildLoadTsFileStatementForSync(dataBaseName, fileAbsolutePath, validateTsFile.get()));
512+
buildLoadTsFileStatementForSync(
513+
dataBaseName,
514+
fileAbsolutePath,
515+
validateTsFile.get(),
516+
shouldConvertDataTypeOnTypeMismatch));
513517
}
514518

515519
static LoadTsFileStatement buildLoadTsFileStatementForSync(
516-
final String dataBaseName, final String fileAbsolutePath, final boolean validateTsFile)
520+
final String dataBaseName,
521+
final String fileAbsolutePath,
522+
final boolean validateTsFile,
523+
final boolean shouldConvertDataTypeOnTypeMismatch)
517524
throws FileNotFoundException {
518525
final LoadTsFileStatement statement = LoadTsFileStatement.createUnchecked(fileAbsolutePath);
519526
statement.setDeleteAfterLoad(true);
520-
statement.setConvertOnTypeMismatch(true);
521-
statement.setVerifySchema(validateTsFile);
527+
statement.setConvertOnTypeMismatch(shouldConvertDataTypeOnTypeMismatch);
528+
statement.setVerifySchema(validateTsFile || shouldConvertDataTypeOnTypeMismatch);
522529
statement.setAutoCreateDatabase(
523530
IoTDBDescriptor.getInstance().getConfig().isAutoCreateSchemaEnabled());
524531
statement.setDatabase(dataBaseName);
@@ -769,9 +776,23 @@ private TSStatus executeStatementWithRetryOnDataTypeMismatch(final Statement sta
769776
TSStatusCode.PIPE_TRANSFER_EXECUTE_STATEMENT_ERROR, "Execute null statement.");
770777
}
771778

779+
// Execute insert statements through the conversion wrapper first to avoid writing a partial
780+
// row/tablet before the type mismatch is converted.
781+
if (shouldConvertDataTypeOnTypeMismatch && statement instanceof InsertBaseStatement) {
782+
final Optional<TSStatus> convertedStatus =
783+
statement.accept(
784+
statementDataTypeConvertExecutionVisitor,
785+
new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
786+
if (convertedStatus.isPresent()) {
787+
return convertedStatus.get();
788+
}
789+
}
790+
772791
final TSStatus status = executeStatement(statement);
773792

774-
// Try to convert the data type and the status code is not success
793+
// Try to convert data type if the status code is not success. Insert statements normally return
794+
// above after the first converted execution. The retry path is kept for load and fallback
795+
// cases.
775796
return shouldConvertDataTypeOnTypeMismatch
776797
&& ((statement instanceof InsertBaseStatement
777798
&& ((InsertBaseStatement) statement).hasFailedMeasurements())

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

Lines changed: 38 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ public void testLoadTsFileSyncStatementUsesTreeDatabaseLevelFromDatabaseName() t
3939
try {
4040
final LoadTsFileStatement statement =
4141
IoTDBDataNodeReceiver.buildLoadTsFileStatementForSync(
42-
"root.test.sg_0", tsFile.toString(), true);
42+
"root.test.sg_0", tsFile.toString(), true, true);
4343

4444
Assert.assertEquals("root.test.sg_0", statement.getDatabase());
4545
Assert.assertEquals(2, statement.getDatabaseLevel());
@@ -54,16 +54,17 @@ public void testLoadTsFileAsyncAttributesUseTreeDatabaseLevelFromDatabaseName()
5454
try {
5555
final Map<String, String> attributes =
5656
IoTDBDataNodeReceiver.buildLoadTsFileAttributesForAsync(
57-
"root.test.sg_0", true, true, true);
57+
"root.test.sg_0", true, false, true);
5858

5959
Assert.assertEquals(
6060
"root.test.sg_0", attributes.get(LoadTsFileConfigurator.DATABASE_NAME_KEY));
6161
Assert.assertEquals("2", attributes.get(LoadTsFileConfigurator.DATABASE_LEVEL_KEY));
6262

6363
final LoadTsFileStatement statement = LoadTsFileStatement.createUnchecked(tsFile.toString());
64-
ActiveLoadPathHelper.applyAttributesToStatement(attributes, statement, true);
64+
ActiveLoadPathHelper.applyAttributesToStatement(attributes, statement, false);
6565
Assert.assertEquals("root.test.sg_0", statement.getDatabase());
6666
Assert.assertEquals(2, statement.getDatabaseLevel());
67+
Assert.assertTrue(statement.isVerifySchema());
6768
} finally {
6869
Files.deleteIfExists(tsFile);
6970
}
@@ -75,7 +76,8 @@ public void testLoadTsFileSyncStatementKeepsDefaultDatabaseLevelWhenDatabaseName
7576
final Path tsFile = Files.createTempFile("pipe-load-default-database-level", ".tsfile");
7677
try {
7778
final LoadTsFileStatement statement =
78-
IoTDBDataNodeReceiver.buildLoadTsFileStatementForSync(null, tsFile.toString(), true);
79+
IoTDBDataNodeReceiver.buildLoadTsFileStatementForSync(
80+
null, tsFile.toString(), true, true);
7981

8082
Assert.assertNull(statement.getDatabase());
8183
Assert.assertEquals(
@@ -92,7 +94,7 @@ public void testRepeatedStatementExceptionLogIsReduced() throws Exception {
9294
try {
9395
final LoadTsFileStatement statement =
9496
IoTDBDataNodeReceiver.buildLoadTsFileStatementForSync(
95-
"root.test.sg_0", tsFile.toString(), true);
97+
"root.test.sg_0", tsFile.toString(), true, true);
9698
final long receiverId = System.nanoTime();
9799
final Exception exception = new RuntimeException("repeated receiver exception " + receiverId);
98100

@@ -107,4 +109,35 @@ public void testRepeatedStatementExceptionLogIsReduced() throws Exception {
107109
Files.deleteIfExists(tsFile);
108110
}
109111
}
112+
113+
@Test
114+
public void testLoadTsFileSyncStatementVerifiesSchemaWhenConvertingType() throws Exception {
115+
final Path tsFile = Files.createTempFile("pipe-load-convert-verify-schema", ".tsfile");
116+
try {
117+
final LoadTsFileStatement statement =
118+
IoTDBDataNodeReceiver.buildLoadTsFileStatementForSync(
119+
"root.test.sg_0", tsFile.toString(), false, true);
120+
121+
Assert.assertTrue(statement.isConvertOnTypeMismatch());
122+
Assert.assertTrue(statement.isVerifySchema());
123+
} finally {
124+
Files.deleteIfExists(tsFile);
125+
}
126+
}
127+
128+
@Test
129+
public void testLoadTsFileSyncStatementCanSkipVerifySchemaWhenNotConvertingType()
130+
throws Exception {
131+
final Path tsFile = Files.createTempFile("pipe-load-no-convert-no-verify-schema", ".tsfile");
132+
try {
133+
final LoadTsFileStatement statement =
134+
IoTDBDataNodeReceiver.buildLoadTsFileStatementForSync(
135+
"root.test.sg_0", tsFile.toString(), false, false);
136+
137+
Assert.assertFalse(statement.isConvertOnTypeMismatch());
138+
Assert.assertFalse(statement.isVerifySchema());
139+
} finally {
140+
Files.deleteIfExists(tsFile);
141+
}
142+
}
110143
}

0 commit comments

Comments
 (0)