Skip to content

Commit 5de1320

Browse files
committed
Bound ConfigNode snapshot buffers and avoid whole-file heap reads
- Add SnapshotStreamFactory providing adaptive, capped and reusable buffered streams for ConfigNode snapshot files, configured via the new config_node_snapshot_buffer_size_max parameter (default 4MB) - PartitionInfo: replace the fixed 32MB snapshot buffer with the adaptive one, sized from the previous snapshot file size - TemplateTable / CNPhysicalPlanGenerator: memory-map template snapshot files instead of copying them into the heap - Edge distribution: cap snapshot buffers at 256KB
1 parent 4697c42 commit 5de1320

13 files changed

Lines changed: 719 additions & 30 deletions

File tree

‎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
@@ -245,6 +245,12 @@ public class ConfigNodeConfig {
245245
private long configNodeRatisConsensusLogAppenderBufferSize = 16 * 1024 * 1024L;
246246
private long schemaRegionRatisConsensusLogAppenderBufferSize = 16 * 1024 * 1024L;
247247

248+
/**
249+
* Upper bound (in bytes) of the in-memory buffer used when taking/loading ConfigNode snapshot
250+
* files. Snapshot streams size their buffers adaptively up to this cap; 0 disables buffering.
251+
*/
252+
private long configNodeSnapshotBufferSizeMax = 4 * 1024 * 1024L;
253+
248254
/**
249255
* RatisConsensus protocol, trigger a snapshot when ratis_snapshot_trigger_threshold logs are
250256
* written.
@@ -973,6 +979,14 @@ public void setSchemaRegionRatisConsensusLogAppenderBufferSize(
973979
schemaRegionRatisConsensusLogAppenderBufferSize;
974980
}
975981

982+
public long getConfigNodeSnapshotBufferSizeMax() {
983+
return configNodeSnapshotBufferSizeMax;
984+
}
985+
986+
public void setConfigNodeSnapshotBufferSizeMax(long configNodeSnapshotBufferSizeMax) {
987+
this.configNodeSnapshotBufferSizeMax = configNodeSnapshotBufferSizeMax;
988+
}
989+
976990
public long getSchemaRegionRatisSnapshotTriggerThreshold() {
977991
return schemaRegionRatisSnapshotTriggerThreshold;
978992
}

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

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import org.apache.iotdb.commons.log.LoggerPeriodicalLogReducer;
2929
import org.apache.iotdb.commons.pipe.config.PipeDescriptor;
3030
import org.apache.iotdb.commons.schema.SchemaConstant;
31+
import org.apache.iotdb.commons.snapshot.SnapshotStreamFactory;
3132
import org.apache.iotdb.commons.utils.NodeUrlUtils;
3233
import org.apache.iotdb.confignode.i18n.ConfigNodeMessages;
3334
import org.apache.iotdb.confignode.manager.load.balancer.RegionBalancer;
@@ -403,6 +404,13 @@ private void loadProperties(TrimProperties properties) throws BadNodeUrlExceptio
403404

404405
loadRatisConsensusConfig(properties);
405406
loadCQConfig(properties);
407+
408+
conf.setConfigNodeSnapshotBufferSizeMax(
409+
Long.parseLong(
410+
properties.getProperty(
411+
"config_node_snapshot_buffer_size_max",
412+
String.valueOf(conf.getConfigNodeSnapshotBufferSizeMax()))));
413+
SnapshotStreamFactory.setBufferSizeMax(conf.getConfigNodeSnapshotBufferSizeMax());
406414
}
407415

408416
private void loadRatisConsensusConfig(TrimProperties properties) {

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java‎

Lines changed: 24 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@
3030
import org.apache.iotdb.commons.partition.SchemaPartitionTable;
3131
import org.apache.iotdb.commons.schema.table.Audit;
3232
import org.apache.iotdb.commons.snapshot.SnapshotProcessor;
33+
import org.apache.iotdb.commons.snapshot.SnapshotStreamFactory;
3334
import org.apache.iotdb.commons.utils.PathUtils;
3435
import org.apache.iotdb.confignode.consensus.request.read.partition.CountTimeSlotListPlan;
3536
import org.apache.iotdb.confignode.consensus.request.read.partition.GetDataPartitionPlan;
@@ -80,11 +81,11 @@
8081
import org.slf4j.Logger;
8182
import org.slf4j.LoggerFactory;
8283

83-
import java.io.BufferedInputStream;
84-
import java.io.BufferedOutputStream;
8584
import java.io.File;
8685
import java.io.FileOutputStream;
8786
import java.io.IOException;
87+
import java.io.InputStream;
88+
import java.io.OutputStream;
8889
import java.nio.file.Files;
8990
import java.util.ArrayList;
9091
import java.util.BitSet;
@@ -119,8 +120,8 @@ public class PartitionInfo implements SnapshotProcessor {
119120

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

122-
// Allocate 8MB buffer for load snapshot of PartitionInfo
123-
private static final int PARTITION_TABLE_BUFFER_SIZE = 32 * 1024 * 1024;
123+
/** Size of the last written snapshot file, used to size the next snapshot buffer adaptively. */
124+
private volatile long lastSnapshotSize = -1L;
124125

125126
/** For Cluster Partition. */
126127
// For allocating Regions
@@ -993,9 +994,13 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
993994
// snapshot operation.
994995
File tmpFile = new File(snapshotFile.getAbsolutePath() + "-" + UUID.randomUUID());
995996

997+
// The write buffer is sized adaptively from the previous snapshot size and capped by
998+
// config_node_snapshot_buffer_size_max, so a small partition table no longer allocates a
999+
// fixed 32MB buffer per snapshot.
9961000
try (FileOutputStream fileOutputStream = new FileOutputStream(tmpFile);
997-
BufferedOutputStream bufferedOutputStream =
998-
new BufferedOutputStream(fileOutputStream, PARTITION_TABLE_BUFFER_SIZE);
1001+
OutputStream bufferedOutputStream =
1002+
SnapshotStreamFactory.createOutputStream(
1003+
fileOutputStream, lastSnapshotSize > 0 ? lastSnapshotSize : 0);
9991004
TIOStreamTransport tioStreamTransport = new TIOStreamTransport(bufferedOutputStream)) {
10001005
TProtocol protocol = new TBinaryProtocol(tioStreamTransport);
10011006

@@ -1024,7 +1029,11 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
10241029
tioStreamTransport.close();
10251030

10261031
// rename file
1027-
return tmpFile.renameTo(snapshotFile);
1032+
final boolean renamed = tmpFile.renameTo(snapshotFile);
1033+
if (renamed) {
1034+
lastSnapshotSize = snapshotFile.length();
1035+
}
1036+
return renamed;
10281037
} finally {
10291038
// with or without success, delete temporary files anyway
10301039
for (int retry = 0; retry < 5; retry++) {
@@ -1050,9 +1059,12 @@ public void processLoadSnapshot(final File snapshotDir) throws TException, IOExc
10501059
return;
10511060
}
10521061

1053-
try (final BufferedInputStream fileInputStream =
1054-
new BufferedInputStream(
1055-
Files.newInputStream(snapshotFile.toPath()), PARTITION_TABLE_BUFFER_SIZE);
1062+
// The read buffer is sized from the file size and capped by
1063+
// config_node_snapshot_buffer_size_max,
1064+
// so loading a snapshot never allocates more than the configured cap.
1065+
try (final InputStream fileInputStream =
1066+
SnapshotStreamFactory.createInputStream(
1067+
Files.newInputStream(snapshotFile.toPath()), snapshotFile.length());
10561068
final TIOStreamTransport tioStreamTransport = new TIOStreamTransport(fileInputStream)) {
10571069
final TProtocol protocol = new TBinaryProtocol(tioStreamTransport);
10581070
// before restoring a snapshot, clear all old data
@@ -1082,6 +1094,8 @@ public void processLoadSnapshot(final File snapshotDir) throws TException, IOExc
10821094
regionMaintainTaskList.add(task);
10831095
}
10841096
}
1097+
// Remember the loaded file size as the estimate for the next snapshot's write buffer.
1098+
lastSnapshotSize = snapshotFile.length();
10851099
}
10861100

10871101
/**

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/CNPhysicalPlanGenerator.java‎

Lines changed: 15 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,6 @@
4747
import org.apache.iotdb.confignode.persistence.schema.mnode.impl.ConfigTableNode;
4848
import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
4949

50-
import org.apache.tsfile.external.commons.io.IOUtils;
5150
import org.apache.tsfile.utils.Pair;
5251
import org.apache.tsfile.utils.ReadWriteIOUtils;
5352
import org.slf4j.Logger;
@@ -60,8 +59,10 @@
6059
import java.io.IOException;
6160
import java.io.InputStream;
6261
import java.nio.ByteBuffer;
62+
import java.nio.channels.FileChannel;
6363
import java.nio.file.Files;
6464
import java.nio.file.Path;
65+
import java.nio.file.StandardOpenOption;
6566
import java.util.ArrayDeque;
6667
import java.util.ArrayList;
6768
import java.util.Collections;
@@ -89,7 +90,9 @@ public class CNPhysicalPlanGenerator
8990

9091
// File input stream.
9192
private InputStream inputStream = null;
92-
private InputStream templateInputStream = null;
93+
94+
// Template snapshot file, memory-mapped on demand instead of being read into the heap.
95+
private Path templateFile = null;
9396

9497
private static final String STRING_ENCODING = "utf-8";
9598

@@ -126,7 +129,7 @@ public CNPhysicalPlanGenerator(final Path schemaInfoFile, final Path templateFil
126129
inputStream = Files.newInputStream(schemaInfoFile);
127130
// Template file is null for table mTree
128131
if (Objects.nonNull(templateFile)) {
129-
templateInputStream = Files.newInputStream(templateFile);
132+
this.templateFile = templateFile;
130133
}
131134
snapshotFileType = CNSnapshotFileType.SCHEMA;
132135
}
@@ -153,7 +156,7 @@ public boolean hasNext() {
153156
} else if (snapshotFileType == CNSnapshotFileType.TTL) {
154157
generateSetTTLPlan();
155158
} else if (snapshotFileType == CNSnapshotFileType.SCHEMA) {
156-
if (Objects.nonNull(templateInputStream)) {
159+
if (Objects.nonNull(templateFile)) {
157160
generateTemplatePlan();
158161
}
159162
if (latestException != null) {
@@ -171,10 +174,6 @@ public boolean hasNext() {
171174
inputStream.close();
172175
inputStream = null;
173176
}
174-
if (templateInputStream != null) {
175-
templateInputStream.close();
176-
templateInputStream = null;
177-
}
178177
} catch (final IOException ioException) {
179178
latestException = ioException;
180179
}
@@ -513,9 +512,14 @@ private void generateDatabasePhysicalPlan() {
513512
}
514513

515514
private void generateTemplatePlan() {
516-
try (final BufferedInputStream bufferedInputStream =
517-
new BufferedInputStream(templateInputStream)) {
518-
final ByteBuffer byteBuffer = ByteBuffer.wrap(IOUtils.toByteArray(bufferedInputStream));
515+
// The template snapshot file is memory-mapped instead of being read into the heap, so
516+
// generating plans from a large template file no longer temporarily doubles its memory
517+
// footprint.
518+
try (final FileChannel channel = FileChannel.open(templateFile, StandardOpenOption.READ)) {
519+
if (channel.size() == 0) {
520+
return;
521+
}
522+
final ByteBuffer byteBuffer = channel.map(FileChannel.MapMode.READ_ONLY, 0, channel.size());
519523
// Skip id
520524
ReadWriteIOUtils.readInt(byteBuffer);
521525
int size = ReadWriteIOUtils.readInt(byteBuffer);

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/TemplateTable.java‎

Lines changed: 8 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -28,23 +28,21 @@
2828

2929
import org.apache.tsfile.common.conf.TSFileDescriptor;
3030
import org.apache.tsfile.enums.TSDataType;
31-
import org.apache.tsfile.external.commons.io.IOUtils;
3231
import org.apache.tsfile.file.metadata.enums.CompressionType;
3332
import org.apache.tsfile.file.metadata.enums.TSEncoding;
3433
import org.apache.tsfile.utils.ReadWriteIOUtils;
3534
import org.apache.tsfile.write.schema.IMeasurementSchema;
3635
import org.slf4j.Logger;
3736
import org.slf4j.LoggerFactory;
3837

39-
import java.io.BufferedInputStream;
4038
import java.io.BufferedOutputStream;
4139
import java.io.File;
42-
import java.io.FileInputStream;
4340
import java.io.FileOutputStream;
4441
import java.io.IOException;
45-
import java.io.InputStream;
4642
import java.io.OutputStream;
4743
import java.nio.ByteBuffer;
44+
import java.nio.channels.FileChannel;
45+
import java.nio.file.StandardOpenOption;
4846
import java.util.ArrayList;
4947
import java.util.List;
5048
import java.util.Map;
@@ -194,8 +192,7 @@ private void serializeTemplate(Template template, OutputStream outputStream) thr
194192
template.serialize(outputStream);
195193
}
196194

197-
private void deserialize(InputStream inputStream) throws IOException {
198-
ByteBuffer byteBuffer = ByteBuffer.wrap(IOUtils.toByteArray(inputStream));
195+
private void deserialize(ByteBuffer byteBuffer) {
199196
templateIdGenerator.set(ReadWriteIOUtils.readInt(byteBuffer));
200197
int size = ReadWriteIOUtils.readInt(byteBuffer);
201198
while (size > 0) {
@@ -259,12 +256,14 @@ public void processLoadSnapshot(File snapshotDir) throws IOException {
259256
return;
260257
}
261258
templateReadWriteLock.writeLock().lock();
262-
try (FileInputStream fileInputStream = new FileInputStream(snapshotFile);
263-
BufferedInputStream bufferedInputStream = new BufferedInputStream(fileInputStream)) {
259+
// The snapshot file is memory-mapped instead of copied into the heap, so loading a large
260+
// template file no longer temporarily doubles its memory footprint.
261+
try (final FileChannel channel =
262+
FileChannel.open(snapshotFile.toPath(), StandardOpenOption.READ)) {
264263
// Load snapshot of template
265264
this.templateMap.clear();
266265
this.templateIdMap.clear();
267-
deserialize(bufferedInputStream);
266+
deserialize(channel.map(FileChannel.MapMode.READ_ONLY, 0, channel.size()));
268267
} finally {
269268
templateReadWriteLock.writeLock().unlock();
270269
}
Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,40 @@
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.confignode.conf;
21+
22+
import org.apache.iotdb.commons.snapshot.SnapshotStreamFactory;
23+
24+
import org.junit.Assert;
25+
import org.junit.Test;
26+
27+
public class ConfigNodeConfigTest {
28+
29+
@Test
30+
public void testSnapshotBufferSizeMaxDefault() {
31+
final ConfigNodeConfig configNodeConfig = new ConfigNodeConfig();
32+
// The code default must stay aligned with SnapshotStreamFactory's default cap.
33+
Assert.assertEquals(
34+
SnapshotStreamFactory.DEFAULT_BUFFER_SIZE_MAX,
35+
configNodeConfig.getConfigNodeSnapshotBufferSizeMax());
36+
37+
configNodeConfig.setConfigNodeSnapshotBufferSizeMax(256 * 1024L);
38+
Assert.assertEquals(256 * 1024L, configNodeConfig.getConfigNodeSnapshotBufferSizeMax());
39+
}
40+
}

‎iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/CNPhysicalPlanGeneratorTest.java‎

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -391,6 +391,37 @@ public void databaseWithoutTemplateGeneratorTest() throws Exception {
391391
Assert.assertEquals(11, count);
392392
}
393393

394+
@Test
395+
public void emptyTemplateFileGeneratorTest() throws Exception {
396+
setupClusterSchemaInfo();
397+
final TDatabaseSchema tDatabaseSchema = new TDatabaseSchema();
398+
tDatabaseSchema.setName("root.sg");
399+
final DatabaseSchemaPlan databaseSchemaPlan =
400+
new DatabaseSchemaPlan(ConfigPhysicalPlanType.CreateDatabase, tDatabaseSchema);
401+
clusterSchemaInfo.createDatabase(databaseSchemaPlan);
402+
Assert.assertTrue(clusterSchemaInfo.processTakeSnapshot(snapshotDir));
403+
404+
final File schemaInfo =
405+
SystemFileFactory.INSTANCE.getFile(snapshotDir + File.separator + SCHEMA_INFO_FILE_NAME);
406+
// A zero-byte template file must be tolerated: no template plans are generated and no
407+
// exception is reported.
408+
final File emptyTemplate =
409+
SystemFileFactory.INSTANCE.getFile(
410+
snapshotDir + File.separator + "empty_" + TEMPLATE_INFO_FILE_NAME);
411+
Assert.assertTrue(emptyTemplate.createNewFile());
412+
413+
final CNPhysicalPlanGenerator planGenerator =
414+
new CNPhysicalPlanGenerator(schemaInfo.toPath(), emptyTemplate.toPath());
415+
int templateCount = 0;
416+
for (ConfigPhysicalPlan plan : planGenerator) {
417+
if (plan.getType() == ConfigPhysicalPlanType.CreateSchemaTemplate) {
418+
templateCount++;
419+
}
420+
}
421+
planGenerator.checkException();
422+
Assert.assertEquals(0, templateCount);
423+
}
424+
394425
@Test
395426
public void templateGeneratorTest() throws Exception {
396427
setupClusterSchemaInfo();

0 commit comments

Comments
 (0)