Skip to content

Commit 99f76a0

Browse files
committed
fix & add IT
1 parent 633c4f0 commit 99f76a0

16 files changed

Lines changed: 611 additions & 196 deletions

File tree

‎integration-test/src/main/java/org/apache/iotdb/it/env/cluster/config/MppDataNodeConfig.java‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -180,4 +180,10 @@ public DataNodeConfig setDnMultiDirStrategy(String multiDirStrategy) {
180180
setProperty("dn_multi_dir_strategy", multiDirStrategy);
181181
return this;
182182
}
183+
184+
@Override
185+
public DataNodeConfig setTableQueryDeviceEntryBatchSizeInBytes(long batchSizeInBytes) {
186+
setProperty("table_query_device_entry_batch_size_in_bytes", String.valueOf(batchSizeInBytes));
187+
return this;
188+
}
183189
}

‎integration-test/src/main/java/org/apache/iotdb/it/env/remote/config/RemoteDataNodeConfig.java‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,4 +125,9 @@ public DataNodeConfig setDnDataDirs(String dnDataDirs) {
125125
public DataNodeConfig setDnMultiDirStrategy(String multiDirStrategy) {
126126
return this;
127127
}
128+
129+
@Override
130+
public DataNodeConfig setTableQueryDeviceEntryBatchSizeInBytes(long batchSizeInBytes) {
131+
return this;
132+
}
128133
}

‎integration-test/src/main/java/org/apache/iotdb/itbase/env/DataNodeConfig.java‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,4 +65,6 @@ DataNodeConfig setLoadActiveListeningCheckIntervalSeconds(
6565
DataNodeConfig setDnDataDirs(String dnDataDirs);
6666

6767
DataNodeConfig setDnMultiDirStrategy(String multiDirStrategy);
68+
69+
DataNodeConfig setTableQueryDeviceEntryBatchSizeInBytes(long batchSizeInBytes);
6870
}
Lines changed: 143 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,143 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
package org.apache.iotdb.relational.it.query.recent;
21+
22+
import org.apache.iotdb.it.env.EnvFactory;
23+
import org.apache.iotdb.itbase.category.TableClusterIT;
24+
import org.apache.iotdb.itbase.category.TableLocalStandaloneIT;
25+
import org.apache.iotdb.itbase.env.BaseEnv;
26+
27+
import org.junit.AfterClass;
28+
import org.junit.BeforeClass;
29+
import org.junit.Test;
30+
import org.junit.experimental.categories.Category;
31+
32+
import java.sql.Connection;
33+
import java.sql.ResultSet;
34+
import java.sql.Statement;
35+
36+
import static org.junit.Assert.assertEquals;
37+
38+
@Category({TableLocalStandaloneIT.class, TableClusterIT.class})
39+
public class IoTDBDeviceEntrySpillIT {
40+
41+
@BeforeClass
42+
public static void setUp() throws Exception {
43+
EnvFactory.getEnv().getConfig().getDataNodeConfig().setTableQueryDeviceEntryBatchSizeInBytes(1);
44+
EnvFactory.getEnv().initClusterEnvironment();
45+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
46+
Statement statement = connection.createStatement()) {
47+
statement.execute("CREATE DATABASE spill_test");
48+
statement.execute(
49+
"CREATE TABLE spill_test.device_data (tag1 STRING TAG, tag2 STRING TAG, "
50+
+ "value INT32 FIELD)");
51+
statement.execute(
52+
"INSERT INTO spill_test.device_data(tag1, tag2, time, value) "
53+
+ "VALUES ('a', 'x', 1, 10), ('a', 'x', 2, 20), "
54+
+ "('b', 'y', 1, 30), ('c', 'z', 1, 40)");
55+
}
56+
}
57+
58+
@AfterClass
59+
public static void tearDown() throws Exception {
60+
EnvFactory.getEnv().cleanClusterEnvironment();
61+
}
62+
63+
@Test
64+
public void testRawFullTableQueryWithSpill() throws Exception {
65+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
66+
Statement statement = connection.createStatement();
67+
ResultSet resultSet =
68+
statement.executeQuery("SELECT tag1, tag2, value FROM spill_test.device_data")) {
69+
int rowCount = 0;
70+
while (resultSet.next()) {
71+
rowCount++;
72+
}
73+
assertEquals(4, rowCount);
74+
}
75+
}
76+
77+
@Test
78+
public void testRawQueriesWithTimeFilterProjectionFilterLimitAndOrdering() throws Exception {
79+
String[] queries = {
80+
"SELECT * FROM spill_test.device_data WHERE time >= 1 AND time < 3",
81+
"SELECT time, value FROM spill_test.device_data WHERE time >= 1 AND time < 3",
82+
"SELECT tag1, tag2, value FROM spill_test.device_data "
83+
+ "WHERE time >= 1 AND time < 3 AND value > 10",
84+
"SELECT * FROM spill_test.device_data WHERE time >= 1 AND time < 3 LIMIT 2",
85+
"SELECT * FROM spill_test.device_data WHERE time >= 1 AND time < 3 ORDER BY time ASC",
86+
"SELECT * FROM spill_test.device_data WHERE time >= 1 AND time < 3 " + "ORDER BY tag1, time"
87+
};
88+
for (String query : queries) {
89+
assertRowCount(query, query.contains("LIMIT 2") ? 2 : query.contains("value > 10") ? 2 : 4);
90+
}
91+
}
92+
93+
@Test
94+
public void testAggregationQueryWithSpill() throws Exception {
95+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
96+
Statement statement = connection.createStatement();
97+
ResultSet resultSet =
98+
statement.executeQuery("SELECT count(value) FROM spill_test.device_data")) {
99+
assertEquals(true, resultSet.next());
100+
assertEquals(4, resultSet.getLong(1));
101+
assertEquals(false, resultSet.next());
102+
}
103+
}
104+
105+
@Test
106+
public void testGroupedAggregationAcrossSpillSegments() throws Exception {
107+
assertRowCount("SELECT tag1, count(*) FROM spill_test.device_data GROUP BY tag1", 3);
108+
assertRowCount(
109+
"SELECT tag1, tag2, count(*), sum(value) FROM spill_test.device_data "
110+
+ "GROUP BY tag1, tag2",
111+
3);
112+
assertRowCount(
113+
"SELECT date_bin(1s, time), count(*) FROM spill_test.device_data "
114+
+ "GROUP BY date_bin(1s, time)",
115+
2);
116+
}
117+
118+
@Test
119+
public void testOrPredicateDoesNotDuplicateDeviceRows() throws Exception {
120+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
121+
Statement statement = connection.createStatement();
122+
ResultSet resultSet =
123+
statement.executeQuery(
124+
"SELECT count(*) FROM spill_test.device_data "
125+
+ "WHERE tag1 = 'a' OR tag2 = 'x'")) {
126+
assertEquals(true, resultSet.next());
127+
assertEquals(2, resultSet.getLong(1));
128+
assertEquals(false, resultSet.next());
129+
}
130+
}
131+
132+
private void assertRowCount(String sql, int expectedRowCount) throws Exception {
133+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
134+
Statement statement = connection.createStatement();
135+
ResultSet resultSet = statement.executeQuery(sql)) {
136+
int rowCount = 0;
137+
while (resultSet.next()) {
138+
rowCount++;
139+
}
140+
assertEquals(expectedRowCount, rowCount);
141+
}
142+
}
143+
}

‎iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3821,5 +3821,7 @@ private DataNodeQueryMessages() {}
38213821
public static final String
38223822
LOG_FAILED_TO_CLEAN_DEVICEENTRY_DATA_SET_ASYNCHRONOUSLY_QUERYID_ARG_PLANNODEID_ARG_9106C4C5 =
38233823
"Failed to clean DeviceEntry data set asynchronously: queryId=%s, planNodeId=%s";
3824+
public static final String LOG_FAILED_TO_CLEAN_DEVICEENTRY_SPILL_DIRECTORY_QUERYID_ARG_ADF95D63 =
3825+
"Failed to clean DeviceEntry spill directory for query %s";
38243826

38253827
}

‎iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4578,5 +4578,7 @@ private DataNodeQueryMessages() {}
45784578
public static final String
45794579
LOG_FAILED_TO_CLEAN_DEVICEENTRY_DATA_SET_ASYNCHRONOUSLY_QUERYID_ARG_PLANNODEID_ARG_9106C4C5 =
45804580
"异步清理 DeviceEntry 数据集失败:queryId=%s,planNodeId=%s";
4581+
public static final String LOG_FAILED_TO_CLEAN_DEVICEENTRY_SPILL_DIRECTORY_QUERYID_ARG_ADF95D63 =
4582+
"清理 query %s 的 DeviceEntry spill 目录失败";
45814583

45824584
}

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceContext.java‎

Lines changed: 7 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -819,14 +819,13 @@ private <T> T initializeQueryDataSource(
819819
return unfinishedResultSupplier.get();
820820
}
821821
QueryDataSource dataSource =
822-
(QueryDataSource)
823-
dataRegion.query(
824-
sourcePaths,
825-
singleDeviceId,
826-
this,
827-
globalTimeFilter != null ? globalTimeFilter.copy() : null,
828-
timePartitions,
829-
waitForLockTime);
822+
dataRegion.query(
823+
sourcePaths,
824+
singleDeviceId,
825+
this,
826+
globalTimeFilter != null ? globalTimeFilter.copy() : null,
827+
timePartitions,
828+
waitForLockTime);
830829
if (dataSource == null) {
831830
return unfinishedResultSupplier.get();
832831
}

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/QueryExecution.java‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@
4848
import org.apache.iotdb.db.queryengine.plan.planner.plan.FragmentInstance;
4949
import org.apache.iotdb.db.queryengine.plan.planner.plan.LogicalQueryPlan;
5050
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeUtil;
51+
import org.apache.iotdb.db.queryengine.plan.relational.metadata.spill.DeviceEntrySpillManager;
5152
import org.apache.iotdb.db.queryengine.plan.scheduler.IScheduler;
5253
import org.apache.iotdb.db.utils.SetThreadName;
5354
import org.apache.iotdb.mpp.rpc.thrift.TFragmentInstanceId;
@@ -437,6 +438,17 @@ private void releaseResource(Throwable t) {
437438
cleanUpResultHandle();
438439
}
439440
context.releaseExternalTsFileQueryResources();
441+
if (t != null) {
442+
try {
443+
DeviceEntrySpillManager.getInstance().deregisterQuery(context.getQueryId().getId());
444+
} catch (Exception e) {
445+
LOGGER.warn(
446+
DataNodeQueryMessages
447+
.LOG_FAILED_TO_CLEAN_DEVICEENTRY_SPILL_DIRECTORY_QUERYID_ARG_ADF95D63,
448+
context.getQueryId().getId(),
449+
e);
450+
}
451+
}
440452
}
441453

442454
/**
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,71 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
package org.apache.iotdb.db.queryengine.plan.relational.metadata.fetcher;
21+
22+
import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId;
23+
import org.apache.iotdb.db.queryengine.common.MPPQueryContext;
24+
import org.apache.iotdb.db.queryengine.plan.relational.metadata.spill.AbstractDeviceEntryMaterializer;
25+
26+
/** Query-scoped state shared while fetching and materializing DeviceEntry objects. */
27+
public final class DeviceEntryFetchContext {
28+
29+
private final MPPQueryContext queryContext;
30+
private final PlanNodeId planNodeId;
31+
private AbstractDeviceEntryMaterializer materializer;
32+
private boolean mayContainDuplicateDevice;
33+
private boolean containsNonAlignedDevice;
34+
35+
public DeviceEntryFetchContext(final MPPQueryContext queryContext, final PlanNodeId planNodeId) {
36+
this.queryContext = queryContext;
37+
this.planNodeId = planNodeId;
38+
}
39+
40+
public MPPQueryContext getQueryContext() {
41+
return queryContext;
42+
}
43+
44+
public PlanNodeId getPlanNodeId() {
45+
return planNodeId;
46+
}
47+
48+
public AbstractDeviceEntryMaterializer getMaterializer() {
49+
return materializer;
50+
}
51+
52+
public void setMaterializer(final AbstractDeviceEntryMaterializer materializer) {
53+
this.materializer = materializer;
54+
}
55+
56+
public boolean mayContainDuplicateDevice() {
57+
return mayContainDuplicateDevice;
58+
}
59+
60+
public void setMayContainDuplicateDevice(final boolean mayContainDuplicateDevice) {
61+
this.mayContainDuplicateDevice = mayContainDuplicateDevice;
62+
}
63+
64+
public boolean containsNonAlignedDevice() {
65+
return containsNonAlignedDevice;
66+
}
67+
68+
public void markContainsNonAlignedDevice() {
69+
containsNonAlignedDevice = true;
70+
}
71+
}

0 commit comments

Comments
 (0)