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
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.apache.iotdb.commons.client.property.ClientPoolProperty.DefaultProperty;
import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.i18n.CommonMessages;
import org.apache.iotdb.confignode.i18n.ConfigNodeMessages;
import org.apache.iotdb.confignode.manager.load.balancer.RegionBalancer;
import org.apache.iotdb.confignode.manager.load.balancer.router.leader.AbstractLeaderBalancer;
Expand Down Expand Up @@ -245,6 +246,9 @@ public class ConfigNodeConfig {
private long configNodeRatisConsensusLogAppenderBufferSize = 16 * 1024 * 1024L;
private long schemaRegionRatisConsensusLogAppenderBufferSize = 16 * 1024 * 1024L;

/** Max size (in bytes) of the in-memory buffer used when taking/loading ConfigNode snapshots. */
private long configNodeSnapshotBufferSizeMax = 4 * 1024 * 1024L;

/**
* RatisConsensus protocol, trigger a snapshot when ratis_snapshot_trigger_threshold logs are
* written.
Expand Down Expand Up @@ -973,6 +977,22 @@ public void setSchemaRegionRatisConsensusLogAppenderBufferSize(
schemaRegionRatisConsensusLogAppenderBufferSize;
}

public long getConfigNodeSnapshotBufferSizeMax() {
return configNodeSnapshotBufferSizeMax;
}

public void setConfigNodeSnapshotBufferSizeMax(long configNodeSnapshotBufferSizeMax) {
if (configNodeSnapshotBufferSizeMax > Integer.MAX_VALUE) {
throw new IllegalArgumentException(
String.format(
CommonMessages
.EXCEPTION_SNAPSHOT_BUFFER_SIZE_MUST_NOT_EXCEED_ARG_BYTES_BUT_WAS_ARG_D1DA6F7E,
Integer.MAX_VALUE,
configNodeSnapshotBufferSizeMax));
}
this.configNodeSnapshotBufferSizeMax = configNodeSnapshotBufferSizeMax;
}

public long getSchemaRegionRatisSnapshotTriggerThreshold() {
return schemaRegionRatisSnapshotTriggerThreshold;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import org.apache.iotdb.commons.log.LoggerPeriodicalLogReducer;
import org.apache.iotdb.commons.pipe.config.PipeDescriptor;
import org.apache.iotdb.commons.schema.SchemaConstant;
import org.apache.iotdb.commons.snapshot.SnapshotStreamFactory;
import org.apache.iotdb.commons.utils.NodeUrlUtils;
import org.apache.iotdb.confignode.i18n.ConfigNodeMessages;
import org.apache.iotdb.confignode.manager.load.balancer.RegionBalancer;
Expand Down Expand Up @@ -403,6 +404,13 @@ private void loadProperties(TrimProperties properties) throws BadNodeUrlExceptio

loadRatisConsensusConfig(properties);
loadCQConfig(properties);

conf.setConfigNodeSnapshotBufferSizeMax(
Long.parseLong(
properties.getProperty(
"config_node_snapshot_buffer_size_max",
String.valueOf(conf.getConfigNodeSnapshotBufferSizeMax()))));
SnapshotStreamFactory.setBufferSizeMax(conf.getConfigNodeSnapshotBufferSizeMax());
}

private void loadRatisConsensusConfig(TrimProperties properties) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import org.apache.iotdb.commons.partition.SchemaPartitionTable;
import org.apache.iotdb.commons.schema.table.Audit;
import org.apache.iotdb.commons.snapshot.SnapshotProcessor;
import org.apache.iotdb.commons.snapshot.SnapshotStreamFactory;
import org.apache.iotdb.commons.utils.PathUtils;
import org.apache.iotdb.confignode.consensus.request.read.partition.CountTimeSlotListPlan;
import org.apache.iotdb.confignode.consensus.request.read.partition.GetDataPartitionPlan;
Expand Down Expand Up @@ -80,11 +81,11 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.BufferedInputStream;
import java.io.BufferedOutputStream;
import java.io.File;
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.nio.file.Files;
import java.util.ArrayList;
import java.util.BitSet;
Expand Down Expand Up @@ -119,9 +120,6 @@ public class PartitionInfo implements SnapshotProcessor {

private static final Logger LOGGER = LoggerFactory.getLogger(PartitionInfo.class);

// Allocate 8MB buffer for load snapshot of PartitionInfo
private static final int PARTITION_TABLE_BUFFER_SIZE = 32 * 1024 * 1024;

/** For Cluster Partition. */
// For allocating Regions
private final AtomicInteger nextRegionGroupId;
Expand Down Expand Up @@ -993,9 +991,11 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
// snapshot operation.
File tmpFile = new File(snapshotFile.getAbsolutePath() + "-" + UUID.randomUUID());

// The write buffer is bounded by config_node_snapshot_buffer_size_max, so a small partition
// table no longer allocates a fixed 32MB buffer per snapshot.
try (FileOutputStream fileOutputStream = new FileOutputStream(tmpFile);
BufferedOutputStream bufferedOutputStream =
new BufferedOutputStream(fileOutputStream, PARTITION_TABLE_BUFFER_SIZE);
OutputStream bufferedOutputStream =
SnapshotStreamFactory.createOutputStream(fileOutputStream);
TIOStreamTransport tioStreamTransport = new TIOStreamTransport(bufferedOutputStream)) {
TProtocol protocol = new TBinaryProtocol(tioStreamTransport);

Expand Down Expand Up @@ -1050,9 +1050,12 @@ public void processLoadSnapshot(final File snapshotDir) throws TException, IOExc
return;
}

try (final BufferedInputStream fileInputStream =
new BufferedInputStream(
Files.newInputStream(snapshotFile.toPath()), PARTITION_TABLE_BUFFER_SIZE);
// The read buffer is sized from the file size and capped by
// config_node_snapshot_buffer_size_max,
// so loading a snapshot never allocates more than the configured cap.
try (final InputStream fileInputStream =
SnapshotStreamFactory.createInputStream(
Files.newInputStream(snapshotFile.toPath()), snapshotFile.length());
final TIOStreamTransport tioStreamTransport = new TIOStreamTransport(fileInputStream)) {
final TProtocol protocol = new TBinaryProtocol(tioStreamTransport);
// before restoring a snapshot, clear all old data
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
/*
* 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.confignode.conf;

import org.apache.iotdb.commons.snapshot.SnapshotStreamFactory;

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

public class ConfigNodeConfigTest {

@Test
public void testSnapshotBufferSizeMaxDefault() {
final ConfigNodeConfig configNodeConfig = new ConfigNodeConfig();
// The code default must stay aligned with SnapshotStreamFactory's default cap.
Assert.assertEquals(
SnapshotStreamFactory.DEFAULT_BUFFER_SIZE_MAX,
configNodeConfig.getConfigNodeSnapshotBufferSizeMax());

configNodeConfig.setConfigNodeSnapshotBufferSizeMax(256 * 1024L);
Assert.assertEquals(256 * 1024L, configNodeConfig.getConfigNodeSnapshotBufferSizeMax());
Assert.assertThrows(
IllegalArgumentException.class,
() -> configNodeConfig.setConfigNodeSnapshotBufferSizeMax((long) Integer.MAX_VALUE + 1));
Assert.assertEquals(256 * 1024L, configNodeConfig.getConfigNodeSnapshotBufferSizeMax());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
import org.apache.iotdb.commons.partition.DataPartitionTable;
import org.apache.iotdb.commons.partition.SchemaPartitionTable;
import org.apache.iotdb.commons.partition.SeriesPartitionTable;
import org.apache.iotdb.commons.snapshot.SnapshotStreamFactory;
import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlanType;
import org.apache.iotdb.confignode.consensus.request.read.region.GetRegionInfoListPlan;
import org.apache.iotdb.confignode.consensus.request.write.database.DatabaseSchemaPlan;
Expand Down Expand Up @@ -66,6 +67,8 @@ public class PartitionInfoTest {
private static PartitionInfo partitionInfo;
private static final File snapshotDir = new File(BASE_OUTPUT_PATH, "snapshot");

private long originalSnapshotBufferSizeMax;

public enum testFlag {
DataPartition(20),
SchemaPartition(30);
Expand All @@ -87,6 +90,10 @@ public void setup() {
if (!snapshotDir.exists()) {
snapshotDir.mkdirs();
}
// Run the snapshot round-trips of this class with a small buffer cap, proving that snapshot
// correctness does not depend on a large fixed buffer.
originalSnapshotBufferSizeMax = SnapshotStreamFactory.getBufferSizeMax();
SnapshotStreamFactory.setBufferSizeMax(64 * 1024);
}

@After
Expand All @@ -95,6 +102,7 @@ public void cleanup() throws IOException {
if (snapshotDir.exists()) {
FileUtils.deleteDirectory(snapshotDir);
}
SnapshotStreamFactory.setBufferSizeMax(originalSnapshotBufferSizeMax);
}

@Test
Expand Down Expand Up @@ -152,6 +160,63 @@ public void testSnapshot() throws TException, IOException {
Assert.assertEquals(partitionInfo, partitionInfo1);
}

@Test
public void testSnapshotWithWriteBufferSmallerThanSnapshot() throws TException, IOException {
partitionInfo.generateNextRegionGroupId();

// Set StorageGroup
partitionInfo.createDatabase(
new DatabaseSchemaPlan(
ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema("root.test")));

// Create a SchemaRegion
CreateRegionGroupsPlan createRegionGroupsReq = new CreateRegionGroupsPlan();
final TRegionReplicaSet schemaRegionReplicaSet =
generateTRegionReplicaSet(
testFlag.SchemaPartition.getFlag(),
generateTConsensusGroupId(
testFlag.SchemaPartition.getFlag(), TConsensusGroupType.SchemaRegion));
createRegionGroupsReq.addRegionGroup("root.test", schemaRegionReplicaSet);
partitionInfo.createRegionGroups(createRegionGroupsReq);

// Create a DataRegion
createRegionGroupsReq = new CreateRegionGroupsPlan();
final TRegionReplicaSet dataRegionReplicaSet =
generateTRegionReplicaSet(
testFlag.DataPartition.getFlag(),
generateTConsensusGroupId(
testFlag.DataPartition.getFlag(), TConsensusGroupType.DataRegion));
createRegionGroupsReq.addRegionGroup("root.test", dataRegionReplicaSet);
partitionInfo.createRegionGroups(createRegionGroupsReq);

// Create a data partition table far larger than the 64KB write buffer configured in setup(),
// so the buffered snapshot stream must flush and wrap around many times.
final CreateDataPartitionPlan createDataPartitionPlan = new CreateDataPartitionPlan();
final Map<String, DataPartitionTable> dataPartitionMap = new HashMap<>();
final Map<TSeriesPartitionSlot, SeriesPartitionTable> slotInfo = new HashMap<>();
final TConsensusGroupId dataRegionId = dataRegionReplicaSet.getRegionId();
for (int seriesSlot = 0; seriesSlot < 2000; seriesSlot++) {
final Map<TTimePartitionSlot, List<TConsensusGroupId>> relationInfo = new HashMap<>();
for (int timeSlot = 0; timeSlot < 8; timeSlot++) {
relationInfo.put(new TTimePartitionSlot(timeSlot), Collections.singletonList(dataRegionId));
}
slotInfo.put(new TSeriesPartitionSlot(seriesSlot), new SeriesPartitionTable(relationInfo));
}
dataPartitionMap.put("root.test", new DataPartitionTable(slotInfo));
createDataPartitionPlan.setAssignedDataPartition(dataPartitionMap);
partitionInfo.createDataPartition(createDataPartitionPlan);

Assert.assertTrue(partitionInfo.processTakeSnapshot(snapshotDir));

// The snapshot must actually be larger than the 64KB buffer for this test to be meaningful.
final File snapshotFile = new File(snapshotDir, "partition_info.bin");
Assert.assertTrue(snapshotFile.length() > 64 * 1024);

final PartitionInfo partitionInfo1 = new PartitionInfo();
partitionInfo1.processLoadSnapshot(snapshotDir);
Assert.assertEquals(partitionInfo, partitionInfo1);
}

@Test
public void testGetRegionType() {

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,11 @@ schema_region_ratis_preserve_logs_num_when_purge=200
config_node_ratis_periodic_snapshot_interval=1800
schema_region_ratis_periodic_snapshot_interval=1800

# Cap the snapshot I/O buffer of the ConfigNode at 8KB (stock default: 4MB).
# Snapshots are tiny on an edge node, so this keeps their transient heap
# allocation negligible.
config_node_snapshot_buffer_size_max=8192

# ---- realtime pipe sync out of the box ----
# The pipe memory pool is 10% of the heap (~22MB at the default 224M budget),
# while the stock pipe memory estimates are sized for datacenter nodes: each
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2123,6 +2123,11 @@ config_node_ratis_snapshot_trigger_threshold=400000
schema_region_ratis_snapshot_trigger_threshold=400000
data_region_ratis_snapshot_trigger_threshold=400000

# max size (in byte) of the in-memory buffer used when taking/loading ConfigNode snapshot files
# effectiveMode: restart
# Datatype: long
config_node_snapshot_buffer_size_max=4194304

# allow flushing Raft Log asynchronously
# effectiveMode: restart
# Datatype: Boolean
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -342,4 +342,7 @@ private CommonMessages() {}
EXCEPTION_XCORR_REQUIRES_EXACTLY_TWO_CALCULATION_COLUMNS_BUT_FOUND_ARG_2FF8EB0C =
"XCorr requires exactly two calculation columns, but found %d.";
public static final String EXCEPTION_COLUMN_LACK_OF_NAME = "the column in table lack of the name";
public static final String
EXCEPTION_SNAPSHOT_BUFFER_SIZE_MUST_NOT_EXCEED_ARG_BYTES_BUT_WAS_ARG_D1DA6F7E =
"Snapshot buffer size must not exceed %d bytes, but was %d.";
}
Original file line number Diff line number Diff line change
Expand Up @@ -238,4 +238,7 @@ private CommonMessages() {}
EXCEPTION_XCORR_REQUIRES_EXACTLY_TWO_CALCULATION_COLUMNS_BUT_FOUND_ARG_2FF8EB0C =
"XCorr 要求必须正好有两列计算列,但实际找到 %d 列。";
public static final String EXCEPTION_COLUMN_LACK_OF_NAME = "表参数中列缺少名字";
public static final String
EXCEPTION_SNAPSHOT_BUFFER_SIZE_MUST_NOT_EXCEED_ARG_BYTES_BUT_WAS_ARG_D1DA6F7E =
"快照缓冲区大小不得超过 %d 字节,但实际为 %d。";
}
Loading
Loading