Skip to content
Open
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 @@ -148,6 +148,7 @@ public class SchemaRegionMemoryImpl implements ISchemaRegion {

private boolean isRecovering = true;
private volatile boolean initialized = false;
private boolean ratisLogAppenderBufferReserved = false;

private final String storageGroupDirPath;
private final String schemaRegionDirPath;
Expand Down Expand Up @@ -198,7 +199,10 @@ public synchronized void init() throws MetadataException {
return;
}

if (config.getSchemaRegionConsensusProtocolClass().equals(ConsensusFactory.RATIS_CONSENSUS)) {
if (!ratisLogAppenderBufferReserved
&& config
.getSchemaRegionConsensusProtocolClass()
.equals(ConsensusFactory.RATIS_CONSENSUS)) {
final long memCost = config.getSchemaRatisConsensusLogAppenderBufferSizeMax();
if (!SystemInfo.getInstance().addDirectBufferMemoryCost(memCost)) {
throw new MetadataException(
Expand All @@ -207,6 +211,7 @@ public synchronized void init() throws MetadataException {
+ ", which is greater than limit mem cost: "
+ SystemInfo.getInstance().getTotalDirectBufferMemorySizeLimit());
}
ratisLogAppenderBufferReserved = true;
}

initDir();
Expand Down Expand Up @@ -425,9 +430,10 @@ public synchronized void deleteSchemaRegion() throws MetadataException {

// delete all the schema region files
SchemaRegionUtils.deleteSchemaRegionFolder(schemaRegionDirPath, logger);
if (config.getSchemaRegionConsensusProtocolClass().equals(ConsensusFactory.RATIS_CONSENSUS)) {
if (ratisLogAppenderBufferReserved) {
SystemInfo.getInstance()
.decreaseDirectBufferMemoryCost(config.getSchemaRatisConsensusLogAppenderBufferSizeMax());
ratisLogAppenderBufferReserved = false;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,7 @@ public class SchemaRegionPBTreeImpl implements ISchemaRegion {
private boolean isRecovering = true;
private volatile boolean initialized = false;
private boolean isClearing = false;
private boolean ratisLogAppenderBufferReserved = false;

private final String storageGroupDirPath;
private final String schemaRegionDirPath;
Expand Down Expand Up @@ -193,7 +194,10 @@ public synchronized void init() throws MetadataException {
return;
}

if (config.getSchemaRegionConsensusProtocolClass().equals(ConsensusFactory.RATIS_CONSENSUS)) {
if (!ratisLogAppenderBufferReserved
&& config
.getSchemaRegionConsensusProtocolClass()
.equals(ConsensusFactory.RATIS_CONSENSUS)) {
long memCost = config.getSchemaRatisConsensusLogAppenderBufferSizeMax();
if (!SystemInfo.getInstance().addDirectBufferMemoryCost(memCost)) {
throw new MetadataException(
Expand All @@ -202,6 +206,7 @@ public synchronized void init() throws MetadataException {
+ ", which is greater than limit mem cost: "
+ SystemInfo.getInstance().getTotalDirectBufferMemorySizeLimit());
}
ratisLogAppenderBufferReserved = true;
}

initDir();
Expand Down Expand Up @@ -481,9 +486,10 @@ public synchronized void deleteSchemaRegion() throws MetadataException {
// delete all the schema region files
SchemaRegionUtils.deleteSchemaRegionFolder(schemaRegionDirPath, logger);

if (config.getSchemaRegionConsensusProtocolClass().equals(ConsensusFactory.RATIS_CONSENSUS)) {
if (ratisLogAppenderBufferReserved) {
SystemInfo.getInstance()
.decreaseDirectBufferMemoryCost(config.getSchemaRatisConsensusLogAppenderBufferSizeMax());
ratisLogAppenderBufferReserved = false;
}
}

Expand All @@ -496,18 +502,20 @@ public boolean createSnapshot(File snapshotDir) {
return false;
}
logger.info("Start create snapshot of schemaRegion {}", schemaRegionId);
boolean isSuccess = true;
boolean isSuccess;
boolean currentResult;
long startTime = System.currentTimeMillis();

long mtreeSnapshotStartTime = System.currentTimeMillis();
isSuccess = isSuccess && mtree.createSnapshot(snapshotDir);
isSuccess = mtree.createSnapshot(snapshotDir);
logger.info(
"MTree snapshot creation of schemaRegion {} costs {}ms.",
schemaRegionId,
System.currentTimeMillis() - mtreeSnapshotStartTime);

long tagSnapshotStartTime = System.currentTimeMillis();
isSuccess = isSuccess && tagManager.createSnapshot(snapshotDir);
currentResult = tagManager.createSnapshot(snapshotDir);
isSuccess = isSuccess && currentResult;
logger.info(
"Tag snapshot creation of schemaRegion {} costs {}ms.",
schemaRegionId,
Expand All @@ -517,7 +525,9 @@ public boolean createSnapshot(File snapshotDir) {
"Snapshot creation of schemaRegion {} costs {}ms.",
schemaRegionId,
System.currentTimeMillis() - startTime);
logger.info("Successfully create snapshot of schemaRegion {}", schemaRegionId);
if (isSuccess) {
logger.info("Successfully create snapshot of schemaRegion {}", schemaRegionId);
}

return isSuccess;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,12 +26,14 @@
import org.apache.iotdb.consensus.ConsensusFactory;
import org.apache.iotdb.db.conf.IoTDBConfig;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.schemaengine.SchemaEngine;
import org.apache.iotdb.db.schemaengine.schemaregion.ISchemaRegion;
import org.apache.iotdb.db.schemaengine.schemaregion.SchemaRegionPlanType;
import org.apache.iotdb.db.schemaengine.schemaregion.read.resp.info.ISchemaInfo;
import org.apache.iotdb.db.schemaengine.schemaregion.read.resp.info.ITimeSeriesSchemaInfo;
import org.apache.iotdb.db.schemaengine.schemaregion.write.req.SchemaRegionWritePlanFactory;
import org.apache.iotdb.db.schemaengine.template.Template;
import org.apache.iotdb.db.storageengine.rescon.memory.SystemInfo;

import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.file.metadata.enums.CompressionType;
Expand Down Expand Up @@ -192,6 +194,44 @@ public void testEmptySnapshot() throws Exception {
}
}

@Test
public void testLoadSnapshotDoesNotReserveDirectBufferMemoryRepeatedly() throws Exception {
String schemaRegionConsensusProtocolClass = config.getSchemaRegionConsensusProtocolClass();
config.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS);
try {
long directBufferMemoryCostBeforeRegion =
SystemInfo.getInstance().getDirectBufferMemoryCost();
ISchemaRegion schemaRegion = getSchemaRegion("root.sg", 0);
long expectedDirectBufferMemoryCost =
directBufferMemoryCostBeforeRegion
+ config.getSchemaRatisConsensusLogAppenderBufferSizeMax();
Assert.assertEquals(
expectedDirectBufferMemoryCost, SystemInfo.getInstance().getDirectBufferMemoryCost());

try {
// A failed snapshot load re-initializes the region. Repeated initialization must not
// reserve
// the Ratis log appender buffer more than once.
File missingSnapshotDir =
new File(config.getSchemaDir() + File.separator + "non-existent-snapshot");
Assert.assertFalse(missingSnapshotDir.exists());

for (int i = 0; i < 2; i++) {
schemaRegion.loadSnapshot(missingSnapshotDir);
Assert.assertEquals(
expectedDirectBufferMemoryCost, SystemInfo.getInstance().getDirectBufferMemoryCost());
}
} finally {
SchemaEngine.getInstance().deleteSchemaRegion(schemaRegion.getSchemaRegionId());
Assert.assertEquals(
directBufferMemoryCostBeforeRegion,
SystemInfo.getInstance().getDirectBufferMemoryCost());
}
} finally {
config.setSchemaRegionConsensusProtocolClass(schemaRegionConsensusProtocolClass);
}
}

@Test
@Ignore
public void testSnapshotPerformance() throws Exception {
Expand Down
Loading