Skip to content

Commit 09eff22

Browse files
committed
Fix pipe history LoadTsFile schema creation
Allow pipe-generated tree LoadTsFile analysis to create the database and missing timeseries schema while verifying the tsfile, even when the receiver disables global auto schema creation. Keep non-pipe load and insert behavior unchanged. Add a regression test for auto-split history sync with deletion and receiver auto-create disabled.
1 parent e948e49 commit 09eff22

6 files changed

Lines changed: 111 additions & 7 deletions

File tree

‎integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/treemodel/auto/basic/IoTDBPipeReceiverAutoCreateDisabledIT.java‎

Lines changed: 48 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@
3636
import java.sql.ResultSetMetaData;
3737
import java.sql.SQLException;
3838
import java.sql.Statement;
39+
import java.util.Arrays;
3940
import java.util.HashSet;
4041
import java.util.Set;
4142

@@ -61,14 +62,16 @@ protected void setupConfig() {
6162
.getConfig()
6263
.getCommonConfig()
6364
.setDataReplicationFactor(1)
64-
.setSchemaReplicationFactor(1);
65+
.setSchemaReplicationFactor(1)
66+
.setPipeAutoSplitFullEnabled(true);
6567
receiverEnv
6668
.getConfig()
6769
.getCommonConfig()
6870
.setAutoCreateSchemaEnabled(false)
6971
.setDatanodeMemoryProportion("3:3:1:1:1:0")
7072
.setDataReplicationFactor(1)
71-
.setSchemaReplicationFactor(1);
73+
.setSchemaReplicationFactor(1)
74+
.setPipeAutoSplitFullEnabled(true);
7275
}
7376

7477
@Test
@@ -122,6 +125,49 @@ public void testReceiverAutoCreateSchemaDisabledWithSpecialTimeSeries() throws E
122125
}
123126
}
124127

128+
@Test
129+
public void testAutoSplitHistoryTsFileWithDeletionWhenReceiverAutoCreateSchemaDisabled()
130+
throws Exception {
131+
TestUtils.executeNonQueries(
132+
senderEnv,
133+
Arrays.asList(
134+
"create database root.sg",
135+
"create timeseries root.sg.non_aligned.s1 with datatype=INT32",
136+
"create timeseries root.sg.non_aligned.s2 with datatype=DOUBLE",
137+
"create aligned timeseries root.sg.aligned(s1 INT32, s2 DOUBLE)",
138+
"insert into root.sg.non_aligned(time, s1, s2) values(1, 1, 1.0), (2, 2, 2.0), (3, 3, 3.0)",
139+
"insert into root.sg.aligned(time, s1, s2) values(1, 10, 10.0), (2, 20, 20.0), (3, 30, 30.0)",
140+
"flush",
141+
"delete from root.sg.non_aligned.s1 where time > 2",
142+
"delete from root.sg.aligned.* where time > 2",
143+
"flush"));
144+
145+
awaitUntilFlush(senderEnv);
146+
147+
TestUtils.executeNonQuery(
148+
senderEnv,
149+
String.format(
150+
"create pipe test with source ('inclusion'='all', 'source.history.enable'='true', 'source.realtime.mode'='batch') "
151+
+ "with sink ('sink'='iotdb-thrift-sink', 'sink.node-urls'='%s')",
152+
receiverEnv.getDataNodeWrapper(0).getIpAndPortString()));
153+
154+
TestUtils.assertDataEventuallyOnEnv(
155+
receiverEnv,
156+
"select * from root.sg.non_aligned",
157+
"Time,root.sg.non_aligned.s1,root.sg.non_aligned.s2,",
158+
new HashSet<>(Arrays.asList("1,1,1.0,", "2,2,2.0,", "3,null,3.0,")));
159+
TestUtils.assertDataEventuallyOnEnv(
160+
receiverEnv,
161+
"select * from root.sg.aligned",
162+
"Time,root.sg.aligned.s1,root.sg.aligned.s2,",
163+
new HashSet<>(Arrays.asList("1,10,10.0,", "2,20,20.0,")));
164+
TestUtils.assertDataEventuallyOnEnv(
165+
receiverEnv,
166+
"show devices root.sg.aligned",
167+
"Device,IsAligned,Template,TTL(ms),",
168+
new HashSet<>(Arrays.asList("root.sg.aligned,true,null,INF,")));
169+
}
170+
125171
private QueryResult queryForResult(final Statement statement, final String sql)
126172
throws SQLException {
127173
try (final ResultSet resultSet = statement.executeQuery(sql)) {

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

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -191,6 +191,10 @@ protected boolean isConvertOnTypeMismatch() {
191191
return isConvertOnTypeMismatch;
192192
}
193193

194+
protected boolean isGeneratedByPipe() {
195+
return isGeneratedByPipe;
196+
}
197+
194198
public IAnalysis analyzeFileByFile(IAnalysis analysis) {
195199
if (!checkBeforeAnalyzeFileByFile(analysis)) {
196200
return analysis;

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

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -215,12 +215,10 @@ private void doAutoCreateAndVerify()
215215
makeSureNoDuplicatedMeasurementsInDevices();
216216
}
217217

218-
if (loadTsFileAnalyzer.isAutoCreateDatabase()) {
218+
if (loadTsFileAnalyzer.isAutoCreateDatabase() || loadTsFileAnalyzer.isGeneratedByPipe()) {
219219
autoCreateDatabase();
220220
}
221221

222-
// schema fetcher will not auto create if config set
223-
// isAutoCreateSchemaEnabled is false.
224222
final ISchemaTree schemaTree = autoCreateSchema();
225223

226224
if (loadTsFileAnalyzer.isVerifySchema()) {
@@ -432,7 +430,8 @@ private ISchemaTree autoCreateSchema() throws IllegalPathException {
432430
encodingsList,
433431
compressionTypesList,
434432
isAlignedList,
435-
loadTsFileAnalyzer.context);
433+
loadTsFileAnalyzer.context,
434+
loadTsFileAnalyzer.isGeneratedByPipe());
436435
}
437436

438437
private void verifySchema(ISchemaTree schemaTree)

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/ClusterSchemaFetcher.java‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -289,6 +289,27 @@ public ISchemaTree fetchSchemaListWithAutoCreate(
289289
final List<CompressionType[]> compressionTypesList,
290290
final List<Boolean> isAlignedList,
291291
final MPPQueryContext context) {
292+
return fetchSchemaListWithAutoCreate(
293+
devicePathList,
294+
measurementsList,
295+
tsDataTypesList,
296+
encodingsList,
297+
compressionTypesList,
298+
isAlignedList,
299+
context,
300+
false);
301+
}
302+
303+
@Override
304+
public ISchemaTree fetchSchemaListWithAutoCreate(
305+
final List<PartialPath> devicePathList,
306+
final List<String[]> measurementsList,
307+
final List<TSDataType[]> tsDataTypesList,
308+
final List<TSEncoding[]> encodingsList,
309+
final List<CompressionType[]> compressionTypesList,
310+
final List<Boolean> isAlignedList,
311+
final MPPQueryContext context,
312+
final boolean forceAutoCreate) {
292313
// The schema cache R/W and fetch operation must be locked together thus the cache clean
293314
// operation executed by delete timeSeries will be effective.
294315
DataNodeSchemaLockManager.getInstance()
@@ -326,7 +347,7 @@ public ISchemaTree fetchSchemaListWithAutoCreate(
326347
schemaTree.mergeSchemaTree(remoteSchemaTree);
327348
}
328349

329-
if (!config.isAutoCreateSchemaEnabled()) {
350+
if (!forceAutoCreate && !config.isAutoCreateSchemaEnabled()) {
330351
// disable auto-create for non-system series
331352
indexOfDevicesWithMissingMeasurements.removeIf(
332353
i ->

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/ISchemaFetcher.java‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,19 @@ ISchemaTree fetchSchemaListWithAutoCreate(
121121
List<Boolean> aligned,
122122
MPPQueryContext context);
123123

124+
default ISchemaTree fetchSchemaListWithAutoCreate(
125+
List<PartialPath> devicePath,
126+
List<String[]> measurements,
127+
List<TSDataType[]> tsDataTypes,
128+
List<TSEncoding[]> encodings,
129+
List<CompressionType[]> compressionTypes,
130+
List<Boolean> aligned,
131+
MPPQueryContext context,
132+
boolean forceAutoCreate) {
133+
return fetchSchemaListWithAutoCreate(
134+
devicePath, measurements, tsDataTypes, encodings, compressionTypes, aligned, context);
135+
}
136+
124137
Pair<Template, PartialPath> checkTemplateSetInfo(PartialPath devicePath);
125138

126139
Pair<Template, PartialPath> checkTemplateSetAndPreSetInfo(

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -97,4 +97,25 @@ public static ISchemaTree validate(
9797
return schemaFetcher.fetchSchemaListWithAutoCreate(
9898
devicePaths, measurements, dataTypes, encodings, compressionTypes, isAlignedList, context);
9999
}
100+
101+
public static ISchemaTree validate(
102+
ISchemaFetcher schemaFetcher,
103+
List<PartialPath> devicePaths,
104+
List<String[]> measurements,
105+
List<TSDataType[]> dataTypes,
106+
List<TSEncoding[]> encodings,
107+
List<CompressionType[]> compressionTypes,
108+
List<Boolean> isAlignedList,
109+
MPPQueryContext context,
110+
boolean forceAutoCreate) {
111+
return schemaFetcher.fetchSchemaListWithAutoCreate(
112+
devicePaths,
113+
measurements,
114+
dataTypes,
115+
encodings,
116+
compressionTypes,
117+
isAlignedList,
118+
context,
119+
forceAutoCreate);
120+
}
100121
}

0 commit comments

Comments
 (0)