Skip to content

Commit 4401337

Browse files
committed
Size PartitionInfo snapshot buffer from the file size instead of a fixed 32MB
PartitionInfo always allocated a fixed 32MB buffer for snapshot load and save, regardless of the actual partition table size. On edge deployments the default heap is only 224MB and snapshots are taken frequently, so a small partition table can still trigger large, short-lived allocations. Size the write buffer from the last snapshot file size and the read buffer from the current file size, still bounded at 32MB so large partition tables keep the previous behavior.
1 parent 4697c42 commit 4401337

2 files changed

Lines changed: 81 additions & 5 deletions

File tree

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

Lines changed: 40 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -119,8 +119,17 @@ public class PartitionInfo implements SnapshotProcessor {
119119

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

122-
// Allocate 8MB buffer for load snapshot of PartitionInfo
123-
private static final int PARTITION_TABLE_BUFFER_SIZE = 32 * 1024 * 1024;
122+
/** Upper bound of a snapshot stream buffer, preserving the previous 32MB cap for large tables. */
123+
private static final int PARTITION_TABLE_BUFFER_SIZE_MAX = 32 * 1024 * 1024;
124+
125+
/** Fallback buffer size used when no snapshot size estimate is available yet. */
126+
private static final int DEFAULT_SNAPSHOT_BUFFER_SIZE = 8 * 1024;
127+
128+
/**
129+
* Size of the last written or loaded snapshot file, used to size the buffer of the next snapshot
130+
* adaptively instead of always allocating a fixed 32MB buffer.
131+
*/
132+
private volatile long lastSnapshotSize = -1L;
124133

125134
/** For Cluster Partition. */
126135
// For allocating Regions
@@ -978,6 +987,19 @@ public Map<TSeriesPartitionSlot, TConsensusGroupId> getLastDataAllotTable(String
978987
return Collections.emptyMap();
979988
}
980989

990+
/**
991+
* Compute the buffer size for a snapshot file of {@code fileSize} bytes: the actual file size
992+
* when known, bounded by {@link #PARTITION_TABLE_BUFFER_SIZE_MAX} so that snapshot I/O never
993+
* allocates more than 32MB, falling back to {@link #DEFAULT_SNAPSHOT_BUFFER_SIZE} when no size
994+
* estimate is available.
995+
*/
996+
static int getSnapshotBufferSize(final long fileSize) {
997+
if (fileSize <= 0) {
998+
return DEFAULT_SNAPSHOT_BUFFER_SIZE;
999+
}
1000+
return (int) Math.min(fileSize, PARTITION_TABLE_BUFFER_SIZE_MAX);
1001+
}
1002+
9811003
@Override
9821004
public boolean processTakeSnapshot(File snapshotDir) throws TException, IOException {
9831005

@@ -993,9 +1015,11 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
9931015
// snapshot operation.
9941016
File tmpFile = new File(snapshotFile.getAbsolutePath() + "-" + UUID.randomUUID());
9951017

1018+
// Size the write buffer from the last snapshot size, so a small partition table no longer
1019+
// allocates a fixed 32MB buffer per snapshot.
9961020
try (FileOutputStream fileOutputStream = new FileOutputStream(tmpFile);
9971021
BufferedOutputStream bufferedOutputStream =
998-
new BufferedOutputStream(fileOutputStream, PARTITION_TABLE_BUFFER_SIZE);
1022+
new BufferedOutputStream(fileOutputStream, getSnapshotBufferSize(lastSnapshotSize));
9991023
TIOStreamTransport tioStreamTransport = new TIOStreamTransport(bufferedOutputStream)) {
10001024
TProtocol protocol = new TBinaryProtocol(tioStreamTransport);
10011025

@@ -1024,7 +1048,12 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
10241048
tioStreamTransport.close();
10251049

10261050
// rename file
1027-
return tmpFile.renameTo(snapshotFile);
1051+
final boolean renamed = tmpFile.renameTo(snapshotFile);
1052+
if (renamed) {
1053+
// Remember the file size to size the buffer of the next snapshot.
1054+
lastSnapshotSize = snapshotFile.length();
1055+
}
1056+
return renamed;
10281057
} finally {
10291058
// with or without success, delete temporary files anyway
10301059
for (int retry = 0; retry < 5; retry++) {
@@ -1050,9 +1079,12 @@ public void processLoadSnapshot(final File snapshotDir) throws TException, IOExc
10501079
return;
10511080
}
10521081

1082+
// Size the read buffer from the file size, so loading a small snapshot no longer allocates a
1083+
// fixed 32MB buffer.
10531084
try (final BufferedInputStream fileInputStream =
10541085
new BufferedInputStream(
1055-
Files.newInputStream(snapshotFile.toPath()), PARTITION_TABLE_BUFFER_SIZE);
1086+
Files.newInputStream(snapshotFile.toPath()),
1087+
getSnapshotBufferSize(snapshotFile.length()));
10561088
final TIOStreamTransport tioStreamTransport = new TIOStreamTransport(fileInputStream)) {
10571089
final TProtocol protocol = new TBinaryProtocol(tioStreamTransport);
10581090
// before restoring a snapshot, clear all old data
@@ -1082,6 +1114,9 @@ public void processLoadSnapshot(final File snapshotDir) throws TException, IOExc
10821114
regionMaintainTaskList.add(task);
10831115
}
10841116
}
1117+
1118+
// Remember the loaded file size to size the buffer of the next snapshot.
1119+
lastSnapshotSize = snapshotFile.length();
10851120
}
10861121

10871122
/**
Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,41 @@
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.persistence.partition;
21+
22+
import org.junit.Assert;
23+
import org.junit.Test;
24+
25+
public class PartitionInfoBufferSizeTest {
26+
27+
@Test
28+
public void testGetSnapshotBufferSize() {
29+
// No size estimate yet: fall back to the default 8KB buffer.
30+
Assert.assertEquals(PartitionInfo.getSnapshotBufferSize(-1L), 8 * 1024);
31+
Assert.assertEquals(PartitionInfo.getSnapshotBufferSize(0L), 8 * 1024);
32+
33+
// A known file size is used directly while it stays below the 32MB cap.
34+
Assert.assertEquals(PartitionInfo.getSnapshotBufferSize(1L), 1);
35+
Assert.assertEquals(PartitionInfo.getSnapshotBufferSize(4 * 1024L), 4 * 1024);
36+
37+
// At and above the cap the buffer is bounded at 32MB.
38+
Assert.assertEquals(PartitionInfo.getSnapshotBufferSize(32L * 1024 * 1024), 32 * 1024 * 1024);
39+
Assert.assertEquals(PartitionInfo.getSnapshotBufferSize(64L * 1024 * 1024), 32 * 1024 * 1024);
40+
}
41+
}

0 commit comments

Comments
 (0)