Skip to content

Commit dc73469

Browse files
authored
Support spill of devce entry in query process (#18493)
1 parent 2c0c307 commit dc73469

99 files changed

Lines changed: 7429 additions & 789 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎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: 146 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,146 @@
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.it.framework.IoTDBTestRunner;
24+
import org.apache.iotdb.itbase.category.TableClusterIT;
25+
import org.apache.iotdb.itbase.category.TableLocalStandaloneIT;
26+
import org.apache.iotdb.itbase.env.BaseEnv;
27+
28+
import org.junit.AfterClass;
29+
import org.junit.BeforeClass;
30+
import org.junit.Test;
31+
import org.junit.experimental.categories.Category;
32+
import org.junit.runner.RunWith;
33+
34+
import java.sql.Connection;
35+
import java.sql.ResultSet;
36+
import java.sql.Statement;
37+
38+
import static org.junit.Assert.assertEquals;
39+
40+
@RunWith(IoTDBTestRunner.class)
41+
@Category({TableLocalStandaloneIT.class, TableClusterIT.class})
42+
public class IoTDBDeviceEntrySpillIT {
43+
44+
@BeforeClass
45+
public static void setUp() throws Exception {
46+
EnvFactory.getEnv().getConfig().getDataNodeConfig().setTableQueryDeviceEntryBatchSizeInBytes(1);
47+
EnvFactory.getEnv().initClusterEnvironment();
48+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
49+
Statement statement = connection.createStatement()) {
50+
statement.execute("CREATE DATABASE spill_test");
51+
statement.execute(
52+
"CREATE TABLE spill_test.device_data (tag1 STRING TAG, tag2 STRING TAG, "
53+
+ "value INT32 FIELD)");
54+
statement.execute(
55+
"INSERT INTO spill_test.device_data(tag1, tag2, time, value) "
56+
+ "VALUES ('a', 'x', 1, 10), ('a', 'x', 2, 20), "
57+
+ "('b', 'y', 1, 30), ('c', 'z', 1, 40)");
58+
}
59+
}
60+
61+
@AfterClass
62+
public static void tearDown() throws Exception {
63+
EnvFactory.getEnv().cleanClusterEnvironment();
64+
}
65+
66+
@Test
67+
public void testRawFullTableQueryWithSpill() throws Exception {
68+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
69+
Statement statement = connection.createStatement();
70+
ResultSet resultSet =
71+
statement.executeQuery("SELECT tag1, tag2, value FROM spill_test.device_data")) {
72+
int rowCount = 0;
73+
while (resultSet.next()) {
74+
rowCount++;
75+
}
76+
assertEquals(4, rowCount);
77+
}
78+
}
79+
80+
@Test
81+
public void testRawQueriesWithTimeFilterProjectionFilterLimitAndOrdering() throws Exception {
82+
String[] queries = {
83+
"SELECT * FROM spill_test.device_data WHERE time >= 1 AND time < 3",
84+
"SELECT time, value FROM spill_test.device_data WHERE time >= 1 AND time < 3",
85+
"SELECT tag1, tag2, value FROM spill_test.device_data "
86+
+ "WHERE time >= 1 AND time < 3 AND value > 10",
87+
"SELECT * FROM spill_test.device_data WHERE time >= 1 AND time < 3 LIMIT 2",
88+
"SELECT * FROM spill_test.device_data WHERE time >= 1 AND time < 3 ORDER BY time ASC",
89+
"SELECT * FROM spill_test.device_data WHERE time >= 1 AND time < 3 " + "ORDER BY tag1, time"
90+
};
91+
for (String query : queries) {
92+
assertRowCount(query, query.contains("LIMIT 2") ? 2 : query.contains("value > 10") ? 3 : 4);
93+
}
94+
}
95+
96+
@Test
97+
public void testAggregationQueryWithSpill() throws Exception {
98+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
99+
Statement statement = connection.createStatement();
100+
ResultSet resultSet =
101+
statement.executeQuery("SELECT count(value) FROM spill_test.device_data")) {
102+
assertEquals(true, resultSet.next());
103+
assertEquals(4, resultSet.getLong(1));
104+
assertEquals(false, resultSet.next());
105+
}
106+
}
107+
108+
@Test
109+
public void testGroupedAggregationAcrossSpillSegments() throws Exception {
110+
assertRowCount("SELECT tag1, count(*) FROM spill_test.device_data GROUP BY tag1", 3);
111+
assertRowCount(
112+
"SELECT tag1, tag2, count(*), sum(value) FROM spill_test.device_data "
113+
+ "GROUP BY tag1, tag2",
114+
3);
115+
assertRowCount(
116+
"SELECT date_bin(1ms, time), count(*) FROM spill_test.device_data "
117+
+ "GROUP BY date_bin(1ms, time)",
118+
2);
119+
}
120+
121+
@Test
122+
public void testOrPredicateDoesNotDuplicateDeviceRows() throws Exception {
123+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
124+
Statement statement = connection.createStatement();
125+
ResultSet resultSet =
126+
statement.executeQuery(
127+
"SELECT count(*) FROM spill_test.device_data "
128+
+ "WHERE tag1 = 'a' OR tag2 = 'x'")) {
129+
assertEquals(true, resultSet.next());
130+
assertEquals(2, resultSet.getLong(1));
131+
assertEquals(false, resultSet.next());
132+
}
133+
}
134+
135+
private void assertRowCount(String sql, int expectedRowCount) throws Exception {
136+
try (Connection connection = EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
137+
Statement statement = connection.createStatement();
138+
ResultSet resultSet = statement.executeQuery(sql)) {
139+
int rowCount = 0;
140+
while (resultSet.next()) {
141+
rowCount++;
142+
}
143+
assertEquals(expectedRowCount, rowCount);
144+
}
145+
}
146+
}

‎iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -152,6 +152,7 @@ public enum TSStatusCode {
152152
REPEATED_RPC_CALL(723),
153153
CANNOT_READ_TSFILE(724),
154154
COPY_TO_WRITE_ERROR(725),
155+
DEVICE_ENTRY_SPILL_NOT_FOUND(726),
155156

156157
// OBJECT
157158
OBJECT_NOT_EXISTS(740),

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1485,5 +1485,8 @@ private DataNodeMiscMessages() {}
14851485
"Missing value arrays for aligned column index %d during expand";
14861486
public static final String EXCEPTION_MISSING_VALUE_ARRAYS_FOR_ALIGNED_COLUMN_INDEX_ARG_DURING_MARK_NULL_VALUE_2893628E =
14871487
"Missing value arrays for aligned column index %d during mark null value";
1488+
public static final String
1489+
LOG_TABLE_QUERY_DEVICE_ENTRY_BATCH_SIZE_IN_BYTES_ARG_EXCEEDS_DN_THRIFT_MAX_FRAME_SIZE_ARG_USING_ARG_AS_THE_EFFECTIVE_VALUE_2AE1BEDA =
1490+
"table_query_device_entry_batch_size_in_bytes (%d) exceeds the maximum RPC payload (dn_thrift_max_frame_size %d minus 1024 bytes); using %d as the effective value";
14881491

14891492
}

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

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3832,5 +3832,49 @@ private DataNodeQueryMessages() {}
38323832
public static final String EXCEPTION_VISIBLEALIASES_IS_NULL_630B27F1 = "visibleAliases is null";
38333833
public static final String EXCEPTION_HAS_NO_PERMISSION_TO_EXECUTE_ARG_BECAUSE_ONLY_THE_SUPERUSER_CAN_ALTER_HIM_HERSELF_C5902893 =
38343834
"Has no permission to execute %s, because only the superuser can alter him/herself.";
3835+
public static final String
3836+
LOG_FAILED_TO_CLEAN_DEVICEENTRY_DATA_SET_ASYNCHRONOUSLY_QUERYID_ARG_PLANNODEID_ARG_9106C4C5 =
3837+
"Failed to clean DeviceEntry data set asynchronously: queryId=%s, planNodeId=%s";
3838+
public static final String LOG_FAILED_TO_CLEAN_DEVICEENTRY_SPILL_DIRECTORY_FOR_QUERY_ARG_53D9C1FC =
3839+
"Failed to clean DeviceEntry spill directory for query %s";
3840+
public static final String
3841+
LOG_FAILED_TO_CLOSE_DEVICEENTRY_SPILL_WRITER_DURING_CLEANUP_ARG_EA7F7941 =
3842+
"Failed to close DeviceEntry spill writer during cleanup: %s";
3843+
public static final String
3844+
LOG_FAILED_TO_CLOSE_DEVICEENTRY_DATA_SET_AFTER_INDEX_SCAN_FAILURE_ARG_57F04319 =
3845+
"Failed to close DeviceEntry data set after index scan failure: %s";
3846+
public static final String
3847+
EXCEPTION_DEVICEENTRY_SPILL_SEGMENT_UNAVAILABLE_MAY_BE_DUE_TO_TIMEOUT_OR_KILL_ARG_B932D10D =
3848+
"DeviceEntry spill segment is unavailable: %s. The query may have timed out or been killed.";
3849+
public static final String
3850+
EXCEPTION_DEVICEENTRY_SEGMENT_AND_ENTRY_COUNTS_MUST_BE_NON_NEGATIVE_7C6E569C =
3851+
"DeviceEntry segment count and entry count must be non-negative";
3852+
public static final String EXCEPTION_OPEN_CONSUMING_READER_IS_NOT_SUPPORTED_8B2A59A4 =
3853+
"Open consuming reader is not supported";
3854+
public static final String
3855+
EXCEPTION_DEVICEENTRY_MATERIALIZER_THRESHOLD_MUST_BE_POSITIVE_ARG_35CF53FD =
3856+
"DeviceEntry materializer threshold must be positive: %d";
3857+
public static final String
3858+
EXCEPTION_DEVICEENTRY_MATERIALIZATION_MEMORY_LIMIT_MUST_BE_POSITIVE_ARG_2EB63404 =
3859+
"DeviceEntry materialization memory limit must be positive: %d";
3860+
public static final String
3861+
EXCEPTION_DEVICEENTRY_DATA_SET_PATH_ESCAPES_THE_QUERY_DIRECTORY_ARG_394A9840 =
3862+
"DeviceEntry data set path escapes the query directory: %s";
3863+
public static final String EXCEPTION_DEVICEENTRY_SEGMENT_ID_MUST_BE_NON_NEGATIVE_ARG_F7653A57 =
3864+
"DeviceEntry segment ID must be non-negative: %d";
3865+
public static final String EXCEPTION_DEVICEENTRY_SEGMENT_HAS_INCOMPLETE_RECORD_LENGTH_0A282C19 =
3866+
"DeviceEntry segment has incomplete record length";
3867+
public static final String
3868+
EXCEPTION_INVALID_DEVICEENTRY_RECORD_LENGTH_ARG_REMAINING_BYTES_ARG_F1C43B72 =
3869+
"Invalid DeviceEntry record length %d, remaining bytes %d";
3870+
public static final String EXCEPTION_DEVICEENTRY_MATERIALIZER_HAS_ALREADY_FINISHED_5538544E =
3871+
"DeviceEntry materializer has already finished";
3872+
public static final String
3873+
EXCEPTION_NO_MATERIALIZER_IS_AVAILABLE_TO_ENFORCE_THE_DEVICEENTRY_MEMORY_LIMIT_FD0604AF =
3874+
"No materializer is available to enforce the DeviceEntry memory limit";
3875+
public static final String EXCEPTION_NO_MORE_DEVICEENTRY_RECORDS_ARE_AVAILABLE_8D51C199 =
3876+
"No more DeviceEntry records are available";
3877+
public static final String EXCEPTION_ONLY_INMEMORYDEVICEENTRYDATASET_SUPPORTS_GET_INLINE_DEVICE_ENTRIES_07A52CAB =
3878+
"Only InMemoryDeviceEntryDataSet supports get inline device entries";
38353879

38363880
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1465,5 +1465,8 @@ private DataNodeMiscMessages() {}
14651465
"扩容过程中缺少对齐列索引 %d 的值数组";
14661466
public static final String EXCEPTION_MISSING_VALUE_ARRAYS_FOR_ALIGNED_COLUMN_INDEX_ARG_DURING_MARK_NULL_VALUE_2893628E =
14671467
"标记空值过程中缺少对齐列索引 %d 的值数组";
1468+
public static final String
1469+
LOG_TABLE_QUERY_DEVICE_ENTRY_BATCH_SIZE_IN_BYTES_ARG_EXCEEDS_DN_THRIFT_MAX_FRAME_SIZE_ARG_USING_ARG_AS_THE_EFFECTIVE_VALUE_2AE1BEDA =
1470+
"table_query_device_entry_batch_size_in_bytes(%d)超过最大 RPC payload(dn_thrift_max_frame_size %d 减去 1024 字节),将使用 %d 作为生效值";
14681471

14691472
}

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

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4590,5 +4590,49 @@ private DataNodeQueryMessages() {}
45904590
public static final String EXCEPTION_VISIBLEALIASES_IS_NULL_630B27F1 = "visibleAliases 不能为空";
45914591
public static final String EXCEPTION_HAS_NO_PERMISSION_TO_EXECUTE_ARG_BECAUSE_ONLY_THE_SUPERUSER_CAN_ALTER_HIM_HERSELF_C5902893 =
45924592
"无权执行 %s,因为只有超级用户可以修改其自身。";
4593+
public static final String
4594+
LOG_FAILED_TO_CLEAN_DEVICEENTRY_DATA_SET_ASYNCHRONOUSLY_QUERYID_ARG_PLANNODEID_ARG_9106C4C5 =
4595+
"异步清理 DeviceEntry 数据集失败:queryId=%s,planNodeId=%s";
4596+
public static final String LOG_FAILED_TO_CLEAN_DEVICEENTRY_SPILL_DIRECTORY_FOR_QUERY_ARG_53D9C1FC =
4597+
"清理 query %s 的 DeviceEntry spill 目录失败";
4598+
public static final String
4599+
LOG_FAILED_TO_CLOSE_DEVICEENTRY_SPILL_WRITER_DURING_CLEANUP_ARG_EA7F7941 =
4600+
"清理过程中关闭 DeviceEntry spill writer 失败:%s";
4601+
public static final String
4602+
LOG_FAILED_TO_CLOSE_DEVICEENTRY_DATA_SET_AFTER_INDEX_SCAN_FAILURE_ARG_57F04319 =
4603+
"index scan 失败后关闭 DeviceEntry data set 失败:%s";
4604+
public static final String
4605+
EXCEPTION_DEVICEENTRY_SPILL_SEGMENT_UNAVAILABLE_MAY_BE_DUE_TO_TIMEOUT_OR_KILL_ARG_B932D10D =
4606+
"DeviceEntry spill segment 不可用:%s。查询可能因超时或主动 kill 而终止。";
4607+
public static final String
4608+
EXCEPTION_DEVICEENTRY_SEGMENT_AND_ENTRY_COUNTS_MUST_BE_NON_NEGATIVE_7C6E569C =
4609+
"DeviceEntry segment count 和 entry count 不能为负数";
4610+
public static final String EXCEPTION_OPEN_CONSUMING_READER_IS_NOT_SUPPORTED_8B2A59A4 =
4611+
"不支持打开消费型读取器";
4612+
public static final String
4613+
EXCEPTION_DEVICEENTRY_MATERIALIZER_THRESHOLD_MUST_BE_POSITIVE_ARG_35CF53FD =
4614+
"DeviceEntry materializer 阈值必须为正数:%d";
4615+
public static final String
4616+
EXCEPTION_DEVICEENTRY_MATERIALIZATION_MEMORY_LIMIT_MUST_BE_POSITIVE_ARG_2EB63404 =
4617+
"DeviceEntry materialization 内存限制必须为正数:%d";
4618+
public static final String
4619+
EXCEPTION_DEVICEENTRY_DATA_SET_PATH_ESCAPES_THE_QUERY_DIRECTORY_ARG_394A9840 =
4620+
"DeviceEntry data set 路径超出 query 目录:%s";
4621+
public static final String EXCEPTION_DEVICEENTRY_SEGMENT_ID_MUST_BE_NON_NEGATIVE_ARG_F7653A57 =
4622+
"DeviceEntry segment ID 不能为负数:%d";
4623+
public static final String EXCEPTION_DEVICEENTRY_SEGMENT_HAS_INCOMPLETE_RECORD_LENGTH_0A282C19 =
4624+
"DeviceEntry segment 的记录长度不完整";
4625+
public static final String
4626+
EXCEPTION_INVALID_DEVICEENTRY_RECORD_LENGTH_ARG_REMAINING_BYTES_ARG_F1C43B72 =
4627+
"DeviceEntry 记录长度 %d 非法,剩余字节数为 %d";
4628+
public static final String EXCEPTION_DEVICEENTRY_MATERIALIZER_HAS_ALREADY_FINISHED_5538544E =
4629+
"DeviceEntry materializer 已结束";
4630+
public static final String
4631+
EXCEPTION_NO_MATERIALIZER_IS_AVAILABLE_TO_ENFORCE_THE_DEVICEENTRY_MEMORY_LIMIT_FD0604AF =
4632+
"没有可用于执行 DeviceEntry 内存限制的 materializer";
4633+
public static final String EXCEPTION_NO_MORE_DEVICEENTRY_RECORDS_ARE_AVAILABLE_8D51C199 =
4634+
"没有更多可用的 DeviceEntry 记录";
4635+
public static final String EXCEPTION_ONLY_INMEMORYDEVICEENTRYDATASET_SUPPORTS_GET_INLINE_DEVICE_ENTRIES_07A52CAB =
4636+
"只有 InMemoryDeviceEntryDataSet 支持获取内存中的设备条目";
45934637

45944638
}

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -259,6 +259,12 @@ public class IoTDBConfig {
259259
private String queryDir =
260260
IoTDBConstant.DN_DEFAULT_DATA_DIR + File.separator + IoTDBConstant.QUERY_FOLDER_NAME;
261261

262+
/**
263+
* Maximum DeviceEntry bytes kept in memory before a table-query spill, capped by the effective
264+
* Thrift frame size minus 1 KiB reserved for the RPC response envelope.
265+
*/
266+
private long tableQueryDeviceEntryBatchSizeInBytes;
267+
262268
/** External lib directory, stores user-uploaded JAR files */
263269
private String extDir = IoTDBConstant.EXT_FOLDER_NAME;
264270

@@ -1792,6 +1798,14 @@ public void setQueryDir(String queryDir) {
17921798
this.queryDir = queryDir;
17931799
}
17941800

1801+
public long getTableQueryDeviceEntryBatchSizeInBytes() {
1802+
return tableQueryDeviceEntryBatchSizeInBytes;
1803+
}
1804+
1805+
public void setTableQueryDeviceEntryBatchSizeInBytes(long tableQueryDeviceEntryBatchSizeInBytes) {
1806+
this.tableQueryDeviceEntryBatchSizeInBytes = tableQueryDeviceEntryBatchSizeInBytes;
1807+
}
1808+
17951809
public String getRatisDataRegionSnapshotDir() {
17961810
return ratisDataRegionSnapshotDir;
17971811
}

0 commit comments

Comments
 (0)