Skip to content

Commit ef9bb52

Browse files
authored
Fix subscription cache memory configuration (#18761)
* Fix subscription cache memory configuration * Include ASM in DataNode runtime dependencies * Update pom.xml
1 parent f8bfcff commit ef9bb52

7 files changed

Lines changed: 147 additions & 8 deletions

File tree

‎iotdb-core/datanode/pom.xml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -293,7 +293,7 @@
293293
<dependency>
294294
<groupId>org.ow2.asm</groupId>
295295
<artifactId>asm</artifactId>
296-
<scope>test</scope>
296+
<scope>runtime</scope>
297297
</dependency>
298298
</dependencies>
299299
<build>

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/cache/SubscriptionPollResponseCache.java‎

Lines changed: 19 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -105,12 +105,14 @@ public static SubscriptionPollResponseCache getInstance() {
105105
}
106106

107107
private SubscriptionPollResponseCache() {
108-
final long initMemorySizeInBytes =
109-
PipeDataNodeResourceManager.memory().getTotalNonFloatingMemorySizeInBytes() / 5;
108+
final long totalNonFloatingMemorySizeInBytes =
109+
PipeDataNodeResourceManager.memory().getTotalNonFloatingMemorySizeInBytes();
110+
final float memoryUsagePercentage =
111+
SubscriptionConfig.getInstance().getSubscriptionCacheMemoryUsagePercentage();
110112
final long maxMemorySizeInBytes =
111-
(long)
112-
(PipeDataNodeResourceManager.memory().getTotalNonFloatingMemorySizeInBytes()
113-
* SubscriptionConfig.getInstance().getSubscriptionCacheMemoryUsagePercentage());
113+
calculateMaxMemorySizeInBytes(totalNonFloatingMemorySizeInBytes, memoryUsagePercentage);
114+
final long initMemorySizeInBytes =
115+
calculateInitialMemorySizeInBytes(totalNonFloatingMemorySizeInBytes, memoryUsagePercentage);
114116

115117
// properties required by pipe memory control framework
116118
final PipeMemoryBlock allocatedMemoryBlock =
@@ -148,4 +150,16 @@ private SubscriptionPollResponseCache() {
148150
newMemory);
149151
});
150152
}
153+
154+
static long calculateInitialMemorySizeInBytes(
155+
final long totalNonFloatingMemorySizeInBytes, final float memoryUsagePercentage) {
156+
return Math.min(
157+
totalNonFloatingMemorySizeInBytes / 5,
158+
calculateMaxMemorySizeInBytes(totalNonFloatingMemorySizeInBytes, memoryUsagePercentage));
159+
}
160+
161+
static long calculateMaxMemorySizeInBytes(
162+
final long totalNonFloatingMemorySizeInBytes, final float memoryUsagePercentage) {
163+
return (long) (totalNonFloatingMemorySizeInBytes * memoryUsagePercentage);
164+
}
151165
}
Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
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.subscription.event.cache;
21+
22+
import org.junit.Assert;
23+
import org.junit.Test;
24+
25+
public class SubscriptionPollResponseCacheTest {
26+
27+
private static final long TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES = 1_000;
28+
29+
@Test
30+
public void testInitialMemoryDoesNotExceedConfiguredMaximum() {
31+
Assert.assertEquals(
32+
50,
33+
SubscriptionPollResponseCache.calculateInitialMemorySizeInBytes(
34+
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.05F));
35+
Assert.assertEquals(
36+
100,
37+
SubscriptionPollResponseCache.calculateInitialMemorySizeInBytes(
38+
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.1F));
39+
Assert.assertEquals(
40+
200,
41+
SubscriptionPollResponseCache.calculateInitialMemorySizeInBytes(
42+
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.2F));
43+
Assert.assertEquals(
44+
200,
45+
SubscriptionPollResponseCache.calculateInitialMemorySizeInBytes(
46+
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.5F));
47+
}
48+
49+
@Test
50+
public void testMaximumMemoryUsesConfiguredPercentage() {
51+
Assert.assertEquals(
52+
50,
53+
SubscriptionPollResponseCache.calculateMaxMemorySizeInBytes(
54+
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.05F));
55+
Assert.assertEquals(
56+
500,
57+
SubscriptionPollResponseCache.calculateMaxMemorySizeInBytes(
58+
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.5F));
59+
}
60+
}

‎iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/CommonMessages.java‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -315,6 +315,9 @@ private CommonMessages() {}
315315
public static final String
316316
EXCEPTION_DISK_SPACE_WARNING_THRESHOLD_MUST_BE_IN_0_1_BUT_WAS_7B345766 =
317317
"disk_space_warning_threshold must be in [0, 1), but was ";
318+
public static final String
319+
EXCEPTION_SUBSCRIPTION_CACHE_MEMORY_USAGE_PERCENTAGE_MUST_BE_IN_0_1_BUT_WAS_ARG_57FE2C66 =
320+
"subscription_cache_memory_usage_percentage must be in [0, 1], but was %s.";
318321
public static final String EXCEPTION_FILTER_FUNCTION_WPASS_VALIDATION =
319322
"the value of wpass should be in (0, 1)";
320323
public static final String EXCEPTION_NO_CALCULATE_COLUMNS = "No columns could be calculated.";

‎iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/CommonMessages.java‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -216,6 +216,9 @@ private CommonMessages() {}
216216
public static final String EXCEPTION_THE_ORDER_BY_CLAUSE_OF_THE_DATA_ARGUMENT_MUST_CONTAIN_EXACTLY_THE_TIME_COLUMN_SPECIFIED_BY_THE_TIMECOL_ARGUMENT_4375BAE9 = "DATA 参数的 ORDER BY 子句必须仅包含 TIMECOL 参数指定的时间列。";
217217
public static final String EXCEPTION_UNSUPPORTED_M4_VALUE_TYPE_AF0EF286 = "不支持的 M4 值类型:";
218218
public static final String EXCEPTION_DISK_SPACE_WARNING_THRESHOLD_MUST_BE_IN_0_1_BUT_WAS_7B345766 = "disk_space_warning_threshold 必须在 [0, 1) 范围内,但实际为 ";
219+
public static final String
220+
EXCEPTION_SUBSCRIPTION_CACHE_MEMORY_USAGE_PERCENTAGE_MUST_BE_IN_0_1_BUT_WAS_ARG_57FE2C66 =
221+
"subscription_cache_memory_usage_percentage 必须在 [0, 1] 范围内,但实际为 %s。";
219222
public static final String LOG_TRUSTED_CHANNEL_FUNCTION_FAILED_INITIATOR_ARG_TARGET_ARG_E4C28443 =
220223
"可信信道功能失效:发起者=%s,目标端=%s";
221224
public static final String

‎iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -344,9 +344,9 @@ public void loadCommonProps(TrimProperties properties) throws IOException {
344344
loadRetryProperties(properties);
345345
}
346346

347-
private void loadSubscriptionProps(TrimProperties properties) {
347+
private void loadSubscriptionProps(TrimProperties properties) throws IOException {
348348
config.setSubscriptionCacheMemoryUsagePercentage(
349-
Float.parseFloat(
349+
parseSubscriptionCacheMemoryUsagePercentage(
350350
properties.getProperty(
351351
"subscription_cache_memory_usage_percentage",
352352
String.valueOf(config.getSubscriptionCacheMemoryUsagePercentage()))));
@@ -563,6 +563,18 @@ private void loadSubscriptionProps(TrimProperties properties) {
563563
String.valueOf(config.getSubscriptionConsensusIdleSafeTimeBarrierIntervalMs()))));
564564
}
565565

566+
static float parseSubscriptionCacheMemoryUsagePercentage(final String value) throws IOException {
567+
final float percentage = Float.parseFloat(value);
568+
if (!Float.isFinite(percentage) || percentage < 0 || percentage > 1) {
569+
throw new IOException(
570+
String.format(
571+
CommonMessages
572+
.EXCEPTION_SUBSCRIPTION_CACHE_MEMORY_USAGE_PERCENTAGE_MUST_BE_IN_0_1_BUT_WAS_ARG_57FE2C66,
573+
percentage));
574+
}
575+
return percentage;
576+
}
577+
566578
public void loadRetryProperties(TrimProperties properties) throws IOException {
567579
config.setRemoteWriteMaxRetryDurationInMs(
568580
Long.parseLong(
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,47 @@
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.commons.conf;
21+
22+
import org.junit.Assert;
23+
import org.junit.Test;
24+
25+
import java.io.IOException;
26+
27+
public class CommonDescriptorSubscriptionCacheMemoryUsagePercentageTest {
28+
29+
@Test
30+
public void testValidPercentage() throws IOException {
31+
for (final String value : new String[] {"0", "0.05", "0.1", "1"}) {
32+
Assert.assertEquals(
33+
Float.parseFloat(value),
34+
CommonDescriptor.parseSubscriptionCacheMemoryUsagePercentage(value),
35+
0);
36+
}
37+
}
38+
39+
@Test
40+
public void testInvalidPercentage() {
41+
for (final String value : new String[] {"-0.01", "1.01", "NaN", "Infinity", "-Infinity"}) {
42+
Assert.assertThrows(
43+
IOException.class,
44+
() -> CommonDescriptor.parseSubscriptionCacheMemoryUsagePercentage(value));
45+
}
46+
}
47+
}

0 commit comments

Comments
 (0)