Skip to content

Commit 109f543

Browse files
committed
Fix ConfigNode snapshot buffering and recovery
1 parent e0fbbbf commit 109f543

21 files changed

Lines changed: 209 additions & 99 deletions

File tree

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

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
import org.apache.iotdb.commons.client.property.ClientPoolProperty.DefaultProperty;
2727
import org.apache.iotdb.commons.conf.CommonDescriptor;
2828
import org.apache.iotdb.commons.conf.IoTDBConstant;
29+
import org.apache.iotdb.commons.i18n.CommonMessages;
2930
import org.apache.iotdb.confignode.i18n.ConfigNodeMessages;
3031
import org.apache.iotdb.confignode.manager.load.balancer.RegionBalancer;
3132
import org.apache.iotdb.confignode.manager.load.balancer.router.leader.AbstractLeaderBalancer;
@@ -981,6 +982,14 @@ public long getConfigNodeSnapshotBufferSizeMax() {
981982
}
982983

983984
public void setConfigNodeSnapshotBufferSizeMax(long configNodeSnapshotBufferSizeMax) {
985+
if (configNodeSnapshotBufferSizeMax > Integer.MAX_VALUE) {
986+
throw new IllegalArgumentException(
987+
String.format(
988+
CommonMessages
989+
.EXCEPTION_SNAPSHOT_BUFFER_SIZE_MUST_NOT_EXCEED_ARG_BYTES_BUT_WAS_ARG_D1DA6F7E,
990+
Integer.MAX_VALUE,
991+
configNodeSnapshotBufferSizeMax));
992+
}
984993
this.configNodeSnapshotBufferSizeMax = configNodeSnapshotBufferSizeMax;
985994
}
986995

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/externalservice/ExternalServiceInfo.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -277,6 +277,7 @@ public boolean processTakeSnapshot(File snapshotDir) throws IOException {
277277
serializeInfos(bufferedOutputStream);
278278

279279
// fsync
280+
bufferedOutputStream.flush();
280281
fileOutputStream.getFD().sync();
281282

282283
return true;

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -235,6 +235,7 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
235235
try (FileOutputStream fileOutputStream = new FileOutputStream(snapshotFile);
236236
OutputStream outputStream = SnapshotStreamFactory.createOutputStream(fileOutputStream)) {
237237
ttlCache.serialize(outputStream);
238+
outputStream.flush();
238239
fileOutputStream.getFD().sync();
239240
} finally {
240241
lock.writeLock().unlock();

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -263,6 +263,7 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
263263
triggerTable.serializeTriggerTable(bufferedOutputStream);
264264

265265
// fsync
266+
bufferedOutputStream.flush();
266267
fileOutputStream.getFD().sync();
267268
return true;
268269
}

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -230,6 +230,7 @@ public boolean processTakeSnapshot(File snapshotDir) throws IOException {
230230
udfTable.serializeUDFTable(bufferedOutputStream);
231231

232232
// fsync
233+
bufferedOutputStream.flush();
233234
fileOutputStream.getFD().sync();
234235

235236
return true;

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/cq/CQInfo.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -246,6 +246,7 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
246246
SnapshotStreamFactory.createOutputStream(fileOutputStream)) {
247247

248248
serialize(bufferedOutputStream);
249+
bufferedOutputStream.flush();
249250
fileOutputStream.getFD().sync();
250251
return true;
251252
} finally {

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/node/NodeInfo.java‎

Lines changed: 14 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,6 @@
2626
import org.apache.iotdb.common.rpc.thrift.TSStatus;
2727
import org.apache.iotdb.commons.cluster.NodeStatus;
2828
import org.apache.iotdb.commons.conf.CommonDescriptor;
29-
import org.apache.iotdb.commons.snapshot.ByteBufferInputStream;
3029
import org.apache.iotdb.commons.snapshot.SnapshotProcessor;
3130
import org.apache.iotdb.commons.snapshot.SnapshotStreamFactory;
3231
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
@@ -61,7 +60,7 @@
6160
import java.io.IOException;
6261
import java.io.InputStream;
6362
import java.io.OutputStream;
64-
import java.nio.ByteBuffer;
63+
import java.nio.channels.Channels;
6564
import java.nio.channels.FileChannel;
6665
import java.nio.file.StandardOpenOption;
6766
import java.util.ArrayList;
@@ -745,11 +744,9 @@ public void processLoadSnapshot(File snapshotDir) throws IOException, TException
745744
aiNodeInfoReadWriteLock.writeLock().lock();
746745
versionInfoReadWriteLock.writeLock().lock();
747746

748-
// The snapshot file is read into a direct buffer instead of the heap, so loading node_info.bin
749-
// no longer allocates a heap array as large as the whole file.
750747
try (final FileChannel channel =
751748
FileChannel.open(snapshotFile.toPath(), StandardOpenOption.READ);
752-
final InputStream inputStream = readSnapshotIntoByteBufferInputStream(channel);
749+
final InputStream inputStream = Channels.newInputStream(channel);
753750
TIOStreamTransport tioStreamTransport = new TIOStreamTransport(inputStream)) {
754751
TProtocol protocol = new TBinaryProtocol(tioStreamTransport);
755752

@@ -763,7 +760,7 @@ public void processLoadSnapshot(File snapshotDir) throws IOException, TException
763760

764761
// TODO: Compatibility design. Should replace this function to actual deserialization method
765762
// in IoTDB 2.2 / 1.5
766-
tryDeserializeRegisteredAINode(inputStream, protocol);
763+
tryDeserializeRegisteredAINode(snapshotFile, channel);
767764

768765
deserializeBuildInfo(inputStream);
769766

@@ -775,22 +772,6 @@ public void processLoadSnapshot(File snapshotDir) throws IOException, TException
775772
}
776773
}
777774

778-
private static ByteBufferInputStream readSnapshotIntoByteBufferInputStream(
779-
final FileChannel channel) throws IOException {
780-
final long fileSize = channel.size();
781-
if (fileSize == 0) {
782-
return new ByteBufferInputStream(ByteBuffer.allocate(0));
783-
}
784-
final ByteBuffer byteBuffer = ByteBuffer.allocateDirect((int) fileSize);
785-
while (byteBuffer.hasRemaining()) {
786-
if (channel.read(byteBuffer) < 0) {
787-
break;
788-
}
789-
}
790-
byteBuffer.flip();
791-
return new ByteBufferInputStream(byteBuffer);
792-
}
793-
794775
private void deserializeRegisteredConfigNode(InputStream inputStream, TProtocol protocol)
795776
throws IOException, TException {
796777
int size = ReadWriteIOUtils.readInt(inputStream);
@@ -815,15 +796,20 @@ private void deserializeRegisteredDataNode(InputStream inputStream, TProtocol pr
815796
}
816797
}
817798

818-
private void tryDeserializeRegisteredAINode(InputStream inputStream, TProtocol protocol)
799+
private void tryDeserializeRegisteredAINode(File snapshotFile, FileChannel channel)
819800
throws IOException {
820-
try {
821-
// 0 has no meaning here
822-
inputStream.mark(0);
823-
deserializeRegisteredAINode(inputStream, protocol);
801+
final long position = channel.position();
802+
try (final FileChannel probeChannel =
803+
FileChannel.open(snapshotFile.toPath(), StandardOpenOption.READ);
804+
final InputStream probeInputStream = Channels.newInputStream(probeChannel);
805+
final TIOStreamTransport probeTransport = new TIOStreamTransport(probeInputStream)) {
806+
probeChannel.position(position);
807+
deserializeRegisteredAINode(probeInputStream, new TBinaryProtocol(probeTransport));
808+
channel.position(probeChannel.position());
824809
} catch (IOException | TException ignore) {
825810
// Exception happens here means that the data is upgraded from the old version
826-
inputStream.reset();
811+
registeredAINodes.clear();
812+
channel.position(position);
827813
}
828814
}
829815

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipePluginInfo.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -416,6 +416,7 @@ public boolean processTakeSnapshot(final File snapshotDir) throws IOException {
416416
final OutputStream bufferedOutputStream =
417417
SnapshotStreamFactory.createOutputStream(fileOutputStream)) {
418418
pipePluginMetaKeeper.processTakeSnapshot(bufferedOutputStream);
419+
bufferedOutputStream.flush();
419420
fileOutputStream.getFD().sync();
420421
}
421422
return true;

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1328,6 +1328,7 @@ public boolean processTakeSnapshot(final File snapshotDir) throws IOException {
13281328
final OutputStream bufferedOutputStream =
13291329
SnapshotStreamFactory.createOutputStream(fileOutputStream)) {
13301330
pipeMetaKeeper.processTakeSnapshot(bufferedOutputStream);
1331+
bufferedOutputStream.flush();
13311332
fileOutputStream.getFD().sync();
13321333
}
13331334
return true;

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/quota/QuotaInfo.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -164,6 +164,7 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
164164
SnapshotStreamFactory.createOutputStream(fileOutputStream)) {
165165
serializeSpaceQuotaLimit(bufferedOutputStream);
166166
serializeThrottleQuotaLimit(bufferedOutputStream);
167+
bufferedOutputStream.flush();
167168
fileOutputStream.getFD().sync();
168169
} finally {
169170
spaceQuotaReadWriteLock.writeLock().unlock();

0 commit comments

Comments
 (0)