disable instrumenting
diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/EntrySerialization.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/EntrySerialization.java index 88385a3..1c5f402 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/EntrySerialization.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/EntrySerialization.java
@@ -19,6 +19,7 @@ package org.apache.iotdb.consensus.natraft.protocol.log; +import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.utils.Timer.Statistic; import org.apache.iotdb.tsfile.compress.ICompressor; import org.apache.iotdb.tsfile.compress.IUnCompressor; @@ -29,25 +30,34 @@ import java.io.IOException; import java.nio.ByteBuffer; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; public class EntrySerialization { private static final Logger logger = LoggerFactory.getLogger(EntrySerialization.class); + private static final ExecutorService serializationExecutor = + IoTDBThreadPoolFactory.newFixedThreadPool(32, "Raft-LogSerialization"); private volatile byte[] recycledBuffer; - private volatile ByteBuffer preSerializationCache; private volatile ByteBuffer serializationCache; private volatile ByteBuffer compressionCache; private CompressionType compressionType = CompressionType.UNCOMPRESSED; private int uncompressedSize; + private Future<ByteBuffer> serializeFuture; public void preSerialize(Entry entry) { - if (preSerializationCache != null || serializationCache != null) { + if (serializeFuture != null || serializationCache != null) { return; } - long startTime = Statistic.SERIALIZE_ENTRY.getOperationStartTime(); - ByteBuffer byteBuffer = entry.serializeInternal(recycledBuffer); - Statistic.SERIALIZE_ENTRY.calOperationCostTimeFromStart(startTime); - preSerializationCache = byteBuffer; + serializeFuture = + serializationExecutor.submit( + () -> { + long startTime = Statistic.SERIALIZE_ENTRY.getOperationStartTime(); + ByteBuffer byteBuffer = entry.serializeInternal(recycledBuffer); + Statistic.SERIALIZE_ENTRY.calOperationCostTimeFromStart(startTime); + return byteBuffer; + }); } public ByteBuffer serialize(Entry entry) { @@ -55,15 +65,23 @@ if (cache != null) { return cache.slice(); } - if (preSerializationCache != null) { - ByteBuffer slice = preSerializationCache.slice(); + if (serializeFuture != null) { + ByteBuffer slice = null; + try { + slice = serializeFuture.get().slice(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } catch (ExecutionException e) { + throw new RuntimeException(e); + } slice.position(1); slice.putLong(entry.getCurrLogIndex()); slice.putLong(entry.getCurrLogTerm()); slice.putLong(entry.getPrevTerm()); slice.position(0); serializationCache = slice; - preSerializationCache = null; + serializeFuture = null; } else { long startTime = Statistic.SERIALIZE_ENTRY.getOperationStartTime(); ByteBuffer byteBuffer = entry.serializeInternal(recycledBuffer); @@ -80,15 +98,23 @@ return cache.slice(); } - if (preSerializationCache != null) { - ByteBuffer slice = preSerializationCache.slice(); + if (serializeFuture != null) { + ByteBuffer slice = null; + try { + slice = serializeFuture.get().slice(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } catch (ExecutionException e) { + throw new RuntimeException(e); + } slice.position(1); slice.putLong(entry.getCurrLogIndex()); slice.putLong(entry.getCurrLogTerm()); slice.putLong(entry.getPrevTerm()); slice.position(0); serializationCache = slice; - preSerializationCache = null; + serializeFuture = null; } else { long startTime = Statistic.SERIALIZE_ENTRY.getOperationStartTime(); ByteBuffer byteBuffer = entry.serializeInternal(recycledBuffer); @@ -149,16 +175,31 @@ ByteBuffer cache; if ((cache = serializationCache) != null) { return cache.remaining(); - } else if ((cache = preSerializationCache) != null) { + } else if ((serializeFuture) != null) { + try { + cache = serializeFuture.get(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } catch (ExecutionException e) { + throw new RuntimeException(e); + } return cache.remaining(); } return 0; } public void clear() { - if (preSerializationCache != null) { - recycledBuffer = preSerializationCache.array(); - preSerializationCache = null; + if (serializeFuture != null) { + try { + recycledBuffer = serializeFuture.get().array(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } catch (ExecutionException e) { + throw new RuntimeException(e); + } + serializeFuture = null; } if (serializationCache != null) { recycledBuffer = serializationCache.array(); @@ -177,14 +218,6 @@ this.recycledBuffer = recycledBuffer; } - public ByteBuffer getPreSerializationCache() { - return preSerializationCache; - } - - public void setPreSerializationCache(ByteBuffer preSerializationCache) { - this.preSerializationCache = preSerializationCache; - } - public ByteBuffer getSerializationCache() { return serializationCache; }
diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/manager/serialization/SyncLogDequeSerializer.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/manager/serialization/SyncLogDequeSerializer.java index 16942ee..68d2206 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/manager/serialization/SyncLogDequeSerializer.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/manager/serialization/SyncLogDequeSerializer.java
@@ -460,10 +460,6 @@ compressedBuffer.position(0); compressedBuffer.limit(compressedLength); - ByteBuffer tempBuffer = ByteBuffer.allocate(compressingBuffer.position()); - unCompressor.uncompress(compressedBuffer, tempBuffer); - compressedBuffer.position(0); - Statistic.PERSISTENCE_COMPRESSED_SIZE.add(compressedLength); Statistic.PERSISTENCE_COMPRESS_TIME.calOperationCostTimeFromStart(startTime); } catch (IOException e) {
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java index 855b49a..37bd8cd 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
@@ -24,7 +24,9 @@ import org.apache.iotdb.commons.consensus.ConsensusGroupId; import org.apache.iotdb.commons.consensus.DataRegionId; import org.apache.iotdb.consensus.ConsensusFactory; +import org.apache.iotdb.consensus.EmptyStateMachine; import org.apache.iotdb.consensus.IConsensus; +import org.apache.iotdb.consensus.IStateMachine; import org.apache.iotdb.consensus.config.ConsensusConfig; import org.apache.iotdb.consensus.config.IoTConsensusConfig; import org.apache.iotdb.consensus.config.IoTConsensusConfig.RPC; @@ -219,11 +221,14 @@ CONF.getDataRegionConsensusProtocolClass()))); } - private static DataRegionStateMachine createDataRegionStateMachine(ConsensusGroupId gid) { + private static IStateMachine createDataRegionStateMachine(ConsensusGroupId gid) { DataRegion dataRegion = StorageEngine.getInstance().getDataRegion((DataRegionId) gid); if (ConsensusFactory.IOT_CONSENSUS.equals(CONF.getDataRegionConsensusProtocolClass())) { return new IoTConsensusDataRegionStateMachine(dataRegion); } else { + if (CONF.isIgnoreStateMachine()) { + return new EmptyStateMachine(); + } return new DataRegionStateMachine(dataRegion); } }
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/flush/MemTableFlushTask.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/flush/MemTableFlushTask.java index d218087..ffbec0f 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/flush/MemTableFlushTask.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/flush/MemTableFlushTask.java
@@ -76,10 +76,7 @@ private final LinkedBlockingQueue<Task> sortTaskQueue = new LinkedBlockingQueue<>(); private final LinkedBlockingQueue<Task> encodingTaskQueue = new LinkedBlockingQueue<>(); - private final LinkedBlockingQueue<Task> ioTaskQueue = - (SystemInfo.getInstance().isEncodingFasterThanIo()) - ? new LinkedBlockingQueue<>(config.getIoTaskQueueSizeForFlushing()) - : new LinkedBlockingQueue<>(); + private final LinkedBlockingQueue<Task> ioTaskQueue = new LinkedBlockingQueue<>(); private String storageGroup; private String dataRegionId;
diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/Timer.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/Timer.java index 5ec874d..64c93b0 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/Timer.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/Timer.java
@@ -30,7 +30,7 @@ private static final Logger logger = LoggerFactory.getLogger(Timer.class); - public static final boolean ENABLE_INSTRUMENTING = true; + public static final boolean ENABLE_INSTRUMENTING = false; private static final String COORDINATOR = "Coordinator"; private static final String META_GROUP_MEMBER = "Meta group member";