Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion iotdb-core/datanode/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -293,7 +293,7 @@
<dependency>
<groupId>org.ow2.asm</groupId>
<artifactId>asm</artifactId>
<scope>test</scope>
<scope>runtime</scope>
</dependency>
</dependencies>
<build>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -105,12 +105,14 @@ public static SubscriptionPollResponseCache getInstance() {
}

private SubscriptionPollResponseCache() {
final long initMemorySizeInBytes =
PipeDataNodeResourceManager.memory().getTotalNonFloatingMemorySizeInBytes() / 5;
final long totalNonFloatingMemorySizeInBytes =
PipeDataNodeResourceManager.memory().getTotalNonFloatingMemorySizeInBytes();
final float memoryUsagePercentage =
SubscriptionConfig.getInstance().getSubscriptionCacheMemoryUsagePercentage();
final long maxMemorySizeInBytes =
(long)
(PipeDataNodeResourceManager.memory().getTotalNonFloatingMemorySizeInBytes()
* SubscriptionConfig.getInstance().getSubscriptionCacheMemoryUsagePercentage());
calculateMaxMemorySizeInBytes(totalNonFloatingMemorySizeInBytes, memoryUsagePercentage);
final long initMemorySizeInBytes =
calculateInitialMemorySizeInBytes(totalNonFloatingMemorySizeInBytes, memoryUsagePercentage);

// properties required by pipe memory control framework
final PipeMemoryBlock allocatedMemoryBlock =
Expand Down Expand Up @@ -148,4 +150,16 @@ private SubscriptionPollResponseCache() {
newMemory);
});
}

static long calculateInitialMemorySizeInBytes(
final long totalNonFloatingMemorySizeInBytes, final float memoryUsagePercentage) {
return Math.min(
totalNonFloatingMemorySizeInBytes / 5,
calculateMaxMemorySizeInBytes(totalNonFloatingMemorySizeInBytes, memoryUsagePercentage));
}

static long calculateMaxMemorySizeInBytes(
final long totalNonFloatingMemorySizeInBytes, final float memoryUsagePercentage) {
return (long) (totalNonFloatingMemorySizeInBytes * memoryUsagePercentage);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

package org.apache.iotdb.db.subscription.event.cache;

import org.junit.Assert;
import org.junit.Test;

public class SubscriptionPollResponseCacheTest {

private static final long TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES = 1_000;

@Test
public void testInitialMemoryDoesNotExceedConfiguredMaximum() {
Assert.assertEquals(
50,
SubscriptionPollResponseCache.calculateInitialMemorySizeInBytes(
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.05F));
Assert.assertEquals(
100,
SubscriptionPollResponseCache.calculateInitialMemorySizeInBytes(
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.1F));
Assert.assertEquals(
200,
SubscriptionPollResponseCache.calculateInitialMemorySizeInBytes(
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.2F));
Assert.assertEquals(
200,
SubscriptionPollResponseCache.calculateInitialMemorySizeInBytes(
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.5F));
}

@Test
public void testMaximumMemoryUsesConfiguredPercentage() {
Assert.assertEquals(
50,
SubscriptionPollResponseCache.calculateMaxMemorySizeInBytes(
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.05F));
Assert.assertEquals(
500,
SubscriptionPollResponseCache.calculateMaxMemorySizeInBytes(
TOTAL_NON_FLOATING_MEMORY_SIZE_IN_BYTES, 0.5F));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -315,6 +315,9 @@ private CommonMessages() {}
public static final String
EXCEPTION_DISK_SPACE_WARNING_THRESHOLD_MUST_BE_IN_0_1_BUT_WAS_7B345766 =
"disk_space_warning_threshold must be in [0, 1), but was ";
public static final String
EXCEPTION_SUBSCRIPTION_CACHE_MEMORY_USAGE_PERCENTAGE_MUST_BE_IN_0_1_BUT_WAS_ARG_57FE2C66 =
"subscription_cache_memory_usage_percentage must be in [0, 1], but was %s.";
public static final String EXCEPTION_FILTER_FUNCTION_WPASS_VALIDATION =
"the value of wpass should be in (0, 1)";
public static final String EXCEPTION_NO_CALCULATE_COLUMNS = "No columns could be calculated.";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,9 @@ private CommonMessages() {}
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 参数指定的时间列。";
public static final String EXCEPTION_UNSUPPORTED_M4_VALUE_TYPE_AF0EF286 = "不支持的 M4 值类型:";
public static final String EXCEPTION_DISK_SPACE_WARNING_THRESHOLD_MUST_BE_IN_0_1_BUT_WAS_7B345766 = "disk_space_warning_threshold 必须在 [0, 1) 范围内,但实际为 ";
public static final String
EXCEPTION_SUBSCRIPTION_CACHE_MEMORY_USAGE_PERCENTAGE_MUST_BE_IN_0_1_BUT_WAS_ARG_57FE2C66 =
"subscription_cache_memory_usage_percentage 必须在 [0, 1] 范围内,但实际为 %s。";
public static final String LOG_TRUSTED_CHANNEL_FUNCTION_FAILED_INITIATOR_ARG_TARGET_ARG_E4C28443 =
"可信信道功能失效:发起者=%s,目标端=%s";
public static final String
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -344,9 +344,9 @@ public void loadCommonProps(TrimProperties properties) throws IOException {
loadRetryProperties(properties);
}

private void loadSubscriptionProps(TrimProperties properties) {
private void loadSubscriptionProps(TrimProperties properties) throws IOException {
config.setSubscriptionCacheMemoryUsagePercentage(
Float.parseFloat(
parseSubscriptionCacheMemoryUsagePercentage(
properties.getProperty(
"subscription_cache_memory_usage_percentage",
String.valueOf(config.getSubscriptionCacheMemoryUsagePercentage()))));
Expand Down Expand Up @@ -563,6 +563,18 @@ private void loadSubscriptionProps(TrimProperties properties) {
String.valueOf(config.getSubscriptionConsensusIdleSafeTimeBarrierIntervalMs()))));
}

static float parseSubscriptionCacheMemoryUsagePercentage(final String value) throws IOException {
final float percentage = Float.parseFloat(value);
if (!Float.isFinite(percentage) || percentage < 0 || percentage > 1) {
throw new IOException(
String.format(
CommonMessages
.EXCEPTION_SUBSCRIPTION_CACHE_MEMORY_USAGE_PERCENTAGE_MUST_BE_IN_0_1_BUT_WAS_ARG_57FE2C66,
percentage));
}
return percentage;
}

public void loadRetryProperties(TrimProperties properties) throws IOException {
config.setRemoteWriteMaxRetryDurationInMs(
Long.parseLong(
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

package org.apache.iotdb.commons.conf;

import org.junit.Assert;
import org.junit.Test;

import java.io.IOException;

public class CommonDescriptorSubscriptionCacheMemoryUsagePercentageTest {

@Test
public void testValidPercentage() throws IOException {
for (final String value : new String[] {"0", "0.05", "0.1", "1"}) {
Assert.assertEquals(
Float.parseFloat(value),
CommonDescriptor.parseSubscriptionCacheMemoryUsagePercentage(value),
0);
}
}

@Test
public void testInvalidPercentage() {
for (final String value : new String[] {"-0.01", "1.01", "NaN", "Infinity", "-Infinity"}) {
Assert.assertThrows(
IOException.class,
() -> CommonDescriptor.parseSubscriptionCacheMemoryUsagePercentage(value));
}
}
}
Loading