Skip to content

Commit 24b56ba

Browse files
committed
Make Edge query and load-event thread pools configurable (#18584)
(cherry picked from commit 9d91245)
1 parent f538fae commit 24b56ba

12 files changed

Lines changed: 347 additions & 5 deletions

File tree

‎integration-test/src/test/java/org/apache/iotdb/edge/it/IoTDBEdgeBasicIT.java‎

Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,10 @@
5858
import java.util.Map;
5959
import java.util.Properties;
6060
import java.util.Set;
61+
import java.util.concurrent.CountDownLatch;
62+
import java.util.concurrent.ExecutorService;
63+
import java.util.concurrent.Executors;
64+
import java.util.concurrent.Future;
6165
import java.util.concurrent.TimeUnit;
6266
import java.util.stream.Collectors;
6367
import java.util.stream.Stream;
@@ -180,11 +184,89 @@ public void testTableModelReadWrite() throws SQLException {
180184
}
181185
}
182186

187+
@Test
188+
public void testConcurrentTableQueriesWithSmallThreadPools() throws Exception {
189+
try (Connection connection = openTableConnection();
190+
Statement statement = connection.createStatement()) {
191+
statement.execute("CREATE DATABASE edge_it_concurrent");
192+
statement.execute("USE edge_it_concurrent");
193+
statement.execute("CREATE TABLE sensor(device STRING TAG, value INT32 FIELD)");
194+
statement.execute("INSERT INTO sensor(time,device,value) VALUES (1,'d1',42), (2,'d1',84)");
195+
}
196+
197+
ExecutorService executor = Executors.newFixedThreadPool(4);
198+
CountDownLatch ready = new CountDownLatch(4);
199+
CountDownLatch start = new CountDownLatch(1);
200+
List<Future<Void>> queries = new ArrayList<>();
201+
try {
202+
for (int i = 0; i < 4; i++) {
203+
queries.add(
204+
executor.submit(
205+
() -> {
206+
try (Connection connection = openTableConnection();
207+
Statement statement = connection.createStatement()) {
208+
statement.execute("USE edge_it_concurrent");
209+
ready.countDown();
210+
assertTrue(start.await(30, TimeUnit.SECONDS));
211+
for (int iteration = 0; iteration < 20; iteration++) {
212+
try (ResultSet result =
213+
statement.executeQuery("SELECT sum(value) FROM sensor")) {
214+
assertTrue(result.next());
215+
assertEquals(126.0, result.getDouble(1), 0.0);
216+
assertFalse(result.next());
217+
}
218+
try (ResultSet result =
219+
statement.executeQuery(
220+
"SELECT value FROM sensor WHERE device='d1' ORDER BY time")) {
221+
assertTrue(result.next());
222+
assertEquals(42, result.getInt(1));
223+
assertTrue(result.next());
224+
assertEquals(84, result.getInt(1));
225+
assertFalse(result.next());
226+
}
227+
}
228+
}
229+
return null;
230+
}));
231+
}
232+
assertTrue(ready.await(30, TimeUnit.SECONDS));
233+
start.countDown();
234+
for (Future<Void> query : queries) {
235+
query.get(60, TimeUnit.SECONDS);
236+
}
237+
} finally {
238+
start.countDown();
239+
executor.shutdownNow();
240+
assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
241+
}
242+
}
243+
183244
private static Connection openTreeConnection() throws SQLException {
184245
return DriverManager.getConnection(
185246
jdbcUrl(), SessionConfig.DEFAULT_USER, SessionConfig.DEFAULT_PASSWORD);
186247
}
187248

249+
@Test
250+
public void testRatisMetadataConsensus() throws SQLException {
251+
Map<String, String> variables = new LinkedHashMap<>();
252+
try (Connection connection = openTableConnection();
253+
Statement statement = connection.createStatement();
254+
ResultSet result = statement.executeQuery("SHOW VARIABLES")) {
255+
while (result.next()) {
256+
variables.put(result.getString(1), result.getString(2));
257+
}
258+
}
259+
assertEquals(
260+
"org.apache.iotdb.consensus.ratis.RatisConsensus",
261+
variables.get("ConfigNodeConsensusProtocolClass"));
262+
assertEquals(
263+
"org.apache.iotdb.consensus.ratis.RatisConsensus",
264+
variables.get("SchemaRegionConsensusProtocolClass"));
265+
assertEquals(
266+
"org.apache.iotdb.consensus.iot.IoTConsensus",
267+
variables.get("DataRegionConsensusProtocolClass"));
268+
}
269+
188270
private static Connection openTableConnection() throws SQLException {
189271
return DriverManager.getConnection(
190272
jdbcUrl() + "?sql_dialect=table",
@@ -195,6 +277,10 @@ private static Connection openTableConnection() throws SQLException {
195277
@Test
196278
public void testPackagedConfiguration() throws Exception {
197279
assertFalse(PACKAGED_SYSTEM_PROPERTIES.containsKey("model_inference_execution_thread_count"));
280+
assertEdgeProperty("coordinator_read_executor_size", "2");
281+
assertEdgeProperty("coordinator_scheduled_executor_size", "2");
282+
assertEdgeProperty("fragment_instance_notification_thread_count", "2");
283+
assertEdgeProperty("cn_load_statistics_publisher_thread_count", "1");
198284
assertEdgeProperty("candidate_compaction_task_queue_size", "10");
199285
assertEdgeProperty("compaction_max_aligned_series_num_in_one_batch", "2");
200286
assertEdgeProperty("target_compaction_file_size", "33554432");

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConfig.java‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -202,6 +202,9 @@ public class ConfigNodeConfig {
202202
private int procedureCoreWorkerThreadsCount =
203203
Math.max(Runtime.getRuntime().availableProcessors() / 4, 16);
204204

205+
/** Thread pool size for publishing cluster load statistics changes. */
206+
private int loadStatisticsPublisherThreadCount = 5;
207+
205208
/** The heartbeat interval in milliseconds. */
206209
private volatile long heartbeatIntervalInMs = 1000;
207210

@@ -738,6 +741,17 @@ public void setProcedureCoreWorkerThreadsCount(int procedureCoreWorkerThreadsCou
738741
this.procedureCoreWorkerThreadsCount = procedureCoreWorkerThreadsCount;
739742
}
740743

744+
public int getLoadStatisticsPublisherThreadCount() {
745+
return loadStatisticsPublisherThreadCount;
746+
}
747+
748+
public void setLoadStatisticsPublisherThreadCount(int loadStatisticsPublisherThreadCount) {
749+
if (loadStatisticsPublisherThreadCount <= 0) {
750+
throw new IllegalArgumentException(CommonMessages.SIZE_MUST_BE_POSITIVE);
751+
}
752+
this.loadStatisticsPublisherThreadCount = loadStatisticsPublisherThreadCount;
753+
}
754+
741755
public long getHeartbeatIntervalInMs() {
742756
return heartbeatIntervalInMs;
743757
}

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeDescriptor.java‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -402,6 +402,12 @@ private void loadProperties(TrimProperties properties) throws BadNodeUrlExceptio
402402
"procedure_core_worker_thread_count",
403403
String.valueOf(conf.getProcedureCoreWorkerThreadsCount()))));
404404

405+
conf.setLoadStatisticsPublisherThreadCount(
406+
Integer.parseInt(
407+
properties.getProperty(
408+
"cn_load_statistics_publisher_thread_count",
409+
String.valueOf(conf.getLoadStatisticsPublisherThreadCount()))));
410+
405411
loadRatisConsensusConfig(properties);
406412
loadCQConfig(properties);
407413

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/EventService.java‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -82,7 +82,10 @@ public EventService(LoadCache loadCache) {
8282
new AsyncEventBus(
8383
ThreadName.CONFIG_NODE_LOAD_PUBLISHER.getName(),
8484
IoTDBThreadPoolFactory.newFixedThreadPool(
85-
5, ThreadName.CONFIG_NODE_LOAD_PUBLISHER.getName()));
85+
ConfigNodeDescriptor.getInstance()
86+
.getConf()
87+
.getLoadStatisticsPublisherThreadCount(),
88+
ThreadName.CONFIG_NODE_LOAD_PUBLISHER.getName()));
8689
}
8790

8891
public void register(final IClusterStatusSubscriber listener) {
Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,85 @@
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+
package org.apache.iotdb.confignode.conf;
20+
21+
import org.apache.iotdb.commons.conf.CommonConfig;
22+
import org.apache.iotdb.commons.conf.ConfigurationFileUtils;
23+
import org.apache.iotdb.commons.conf.TrimProperties;
24+
25+
import org.junit.Rule;
26+
import org.junit.Test;
27+
import org.junit.rules.TemporaryFolder;
28+
29+
import java.io.File;
30+
import java.lang.reflect.Constructor;
31+
import java.nio.charset.StandardCharsets;
32+
import java.nio.file.Files;
33+
34+
import static org.junit.Assert.assertEquals;
35+
import static org.junit.Assert.assertThrows;
36+
37+
public class LoadStatisticsPublisherConfigTest {
38+
39+
@Rule public TemporaryFolder temporaryFolder = new TemporaryFolder();
40+
41+
@Test
42+
public void testDefaultsAndPositiveSize() throws Exception {
43+
ConfigNodeConfig config = new ConfigNodeConfig();
44+
assertEquals(5, config.getLoadStatisticsPublisherThreadCount());
45+
assertEquals(
46+
"5",
47+
ConfigurationFileUtils.getConfigurationDefaultValue(
48+
"cn_load_statistics_publisher_thread_count"));
49+
for (int invalid : new int[] {0, -1}) {
50+
assertThrows(
51+
IllegalArgumentException.class,
52+
() -> config.setLoadStatisticsPublisherThreadCount(invalid));
53+
}
54+
assertEquals(5, config.getLoadStatisticsPublisherThreadCount());
55+
}
56+
57+
@Test
58+
public void testStartupOverrideIsRestartOnly() throws Exception {
59+
String originalConf = System.getProperty(ConfigNodeConstant.CONFIGNODE_CONF);
60+
File confDir = temporaryFolder.newFolder();
61+
Files.writeString(
62+
confDir.toPath().resolve(CommonConfig.SYSTEM_CONFIG_NAME),
63+
"cn_seed_config_node=127.0.0.1:10710\ncn_load_statistics_publisher_thread_count=2\n",
64+
StandardCharsets.UTF_8);
65+
System.setProperty(ConfigNodeConstant.CONFIGNODE_CONF, confDir.getAbsolutePath());
66+
try {
67+
Constructor<ConfigNodeDescriptor> constructor =
68+
ConfigNodeDescriptor.class.getDeclaredConstructor();
69+
constructor.setAccessible(true);
70+
ConfigNodeDescriptor descriptor = constructor.newInstance();
71+
assertEquals(2, descriptor.getConf().getLoadStatisticsPublisherThreadCount());
72+
73+
TrimProperties properties = new TrimProperties();
74+
properties.setProperty("cn_load_statistics_publisher_thread_count", "3");
75+
descriptor.loadHotModifiedProps(properties);
76+
assertEquals(2, descriptor.getConf().getLoadStatisticsPublisherThreadCount());
77+
} finally {
78+
if (originalConf == null) {
79+
System.clearProperty(ConfigNodeConstant.CONFIGNODE_CONF);
80+
} else {
81+
System.setProperty(ConfigNodeConstant.CONFIGNODE_CONF, originalConf);
82+
}
83+
}
84+
}
85+
}

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

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import org.apache.iotdb.commons.conf.CommonDescriptor;
2626
import org.apache.iotdb.commons.conf.IoTDBConstant;
2727
import org.apache.iotdb.commons.enums.ReadConsistencyLevel;
28+
import org.apache.iotdb.commons.i18n.CommonMessages;
2829
import org.apache.iotdb.commons.pipe.config.PipeConfig;
2930
import org.apache.iotdb.commons.utils.FileUtils;
3031
import org.apache.iotdb.consensus.ConsensusFactory;
@@ -1002,6 +1003,12 @@ public class IoTDBConfig {
10021003
/** ThreadPool size for read operation in coordinator */
10031004
private int coordinatorReadExecutorSize = 20;
10041005

1006+
/** Thread pool size for scheduling query state checks and termination. */
1007+
private int coordinatorScheduledExecutorSize = 10;
1008+
1009+
/** Thread pool size for fragment instance state change notifications. */
1010+
private int fragmentInstanceNotificationThreadCount = 4;
1011+
10051012
/** Policy of DataNodeSchemaCache eviction */
10061013
private String dataNodeSchemaCacheEvictionPolicy = "FIFO";
10071014

@@ -3479,9 +3486,35 @@ public int getCoordinatorReadExecutorSize() {
34793486
}
34803487

34813488
public void setCoordinatorReadExecutorSize(int coordinatorReadExecutorSize) {
3489+
if (coordinatorReadExecutorSize <= 0) {
3490+
throw new IllegalArgumentException(CommonMessages.SIZE_MUST_BE_POSITIVE);
3491+
}
34823492
this.coordinatorReadExecutorSize = coordinatorReadExecutorSize;
34833493
}
34843494

3495+
public int getCoordinatorScheduledExecutorSize() {
3496+
return coordinatorScheduledExecutorSize;
3497+
}
3498+
3499+
public void setCoordinatorScheduledExecutorSize(int coordinatorScheduledExecutorSize) {
3500+
if (coordinatorScheduledExecutorSize <= 0) {
3501+
throw new IllegalArgumentException(CommonMessages.SIZE_MUST_BE_POSITIVE);
3502+
}
3503+
this.coordinatorScheduledExecutorSize = coordinatorScheduledExecutorSize;
3504+
}
3505+
3506+
public int getFragmentInstanceNotificationThreadCount() {
3507+
return fragmentInstanceNotificationThreadCount;
3508+
}
3509+
3510+
public void setFragmentInstanceNotificationThreadCount(
3511+
int fragmentInstanceNotificationThreadCount) {
3512+
if (fragmentInstanceNotificationThreadCount <= 0) {
3513+
throw new IllegalArgumentException(CommonMessages.SIZE_MUST_BE_POSITIVE);
3514+
}
3515+
this.fragmentInstanceNotificationThreadCount = fragmentInstanceNotificationThreadCount;
3516+
}
3517+
34853518
public TEndPoint getAddressAndPort() {
34863519
return new TEndPoint(rpcAddress, rpcPort);
34873520
}

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

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -988,6 +988,16 @@ public void loadProperties(TrimProperties properties) throws BadNodeUrlException
988988
properties.getProperty(
989989
"coordinator_read_executor_size",
990990
Integer.toString(conf.getCoordinatorReadExecutorSize()))));
991+
conf.setCoordinatorScheduledExecutorSize(
992+
Integer.parseInt(
993+
properties.getProperty(
994+
"coordinator_scheduled_executor_size",
995+
Integer.toString(conf.getCoordinatorScheduledExecutorSize()))));
996+
conf.setFragmentInstanceNotificationThreadCount(
997+
Integer.parseInt(
998+
properties.getProperty(
999+
"fragment_instance_notification_thread_count",
1000+
Integer.toString(conf.getFragmentInstanceNotificationThreadCount()))));
9911001
conf.setDataNodeTableSchemaCacheSize(
9921002
Long.parseLong(
9931003
properties.getProperty(

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -110,7 +110,8 @@ private FragmentInstanceManager() {
110110
1, ThreadName.FRAGMENT_INSTANCE_MANAGEMENT.getName());
111111
this.instanceNotificationExecutor =
112112
IoTDBThreadPoolFactory.newFixedThreadPool(
113-
4, ThreadName.FRAGMENT_INSTANCE_NOTIFICATION.getName());
113+
IoTDBDescriptor.getInstance().getConfig().getFragmentInstanceNotificationThreadCount(),
114+
ThreadName.FRAGMENT_INSTANCE_NOTIFICATION.getName());
114115

115116
this.infoCacheTime = new Duration(5, TimeUnit.MINUTES);
116117

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -200,7 +200,6 @@
200200
public class Coordinator {
201201

202202
private static final Logger LOGGER = LoggerFactory.getLogger(Coordinator.class);
203-
private static final int COORDINATOR_SCHEDULED_EXECUTOR_SIZE = 10;
204203
private static final IoTDBConfig CONFIG = IoTDBDescriptor.getInstance().getConfig();
205204
private static final CommonConfig COMMON_CONFIG = CommonDescriptor.getInstance().getConfig();
206205

@@ -856,7 +855,7 @@ private ExecutorService getQueryExecutor() {
856855

857856
private ScheduledExecutorService getScheduledExecutor() {
858857
return IoTDBThreadPoolFactory.newScheduledThreadPool(
859-
COORDINATOR_SCHEDULED_EXECUTOR_SIZE,
858+
CONFIG.getCoordinatorScheduledExecutorSize(),
860859
ThreadName.MPP_COORDINATOR_SCHEDULED_EXECUTOR.getName());
861860
}
862861

0 commit comments

Comments
 (0)