This is an automated email from the ASF dual-hosted git repository. rong pushed a commit to branch pipe-hybrid-opt in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit ac2dae556278f9aece1d3bfc56a27aaf237a0de1 Author: Steve Yurong Su <[email protected]> AuthorDate: Wed Oct 25 12:42:27 2023 +0800 Pipe: fine tune hybrid mode by removing tooManyWALPinned judgement and increasing pipeMaxAllowedPendingTsFileEpochPerDataRegion to 2 --- .../event/common/heartbeat/PipeHeartbeatEvent.java | 4 ++-- .../PipeRealtimeDataRegionHybridExtractor.java | 21 ++++++--------------- .../org/apache/iotdb/commons/conf/CommonConfig.java | 2 +- 3 files changed, 9 insertions(+), 18 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEvent.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEvent.java index 397babdc431..9c1ad160f39 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEvent.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEvent.java @@ -174,13 +174,13 @@ public class PipeHeartbeatEvent extends EnrichedEvent { public void recordBufferQueueSize(EnrichedDeque<Event> bufferQueue) { if (shouldPrintMessage) { bufferQueueTabletSize = bufferQueue.getTabletInsertionEventCount(); + bufferQueueTsFileSize = bufferQueue.getTsFileInsertionEventCount(); bufferQueueSize = bufferQueue.size(); } - bufferQueueTsFileSize = bufferQueue.getTsFileInsertionEventCount(); if (extractor instanceof PipeRealtimeDataRegionHybridExtractor) { ((PipeRealtimeDataRegionHybridExtractor) extractor) - .informEventCollectorQueueTsFileSize(bufferQueueTsFileSize); + .informEventCollectorQueueTsFileSize(bufferQueue.getTsFileInsertionEventCount()); } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionHybridExtractor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionHybridExtractor.java index 257479a1d62..7b6e5e42dcf 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionHybridExtractor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionHybridExtractor.java @@ -26,7 +26,6 @@ import org.apache.iotdb.db.pipe.agent.PipeAgent; import org.apache.iotdb.db.pipe.event.common.heartbeat.PipeHeartbeatEvent; import org.apache.iotdb.db.pipe.event.realtime.PipeRealtimeEvent; import org.apache.iotdb.db.pipe.extractor.realtime.epoch.TsFileEpoch; -import org.apache.iotdb.db.pipe.resource.PipeResourceManager; import org.apache.iotdb.db.storageengine.dataregion.wal.WALManager; import org.apache.iotdb.pipe.api.event.Event; import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent; @@ -76,15 +75,13 @@ public class PipeRealtimeDataRegionHybridExtractor extends PipeRealtimeDataRegio private void extractTabletInsertion(PipeRealtimeEvent event) { if (!isStartedToSupply || mayWalSizeReachThrottleThreshold() - || tooManyWALPinned() || isTsFileEventCountInQueueExceededLimit()) { // In the following 3 cases, we should not extract any more tablet events. all the data // represented by the tablet events should be carried by the following tsfile event: // 1. The historical extractor has not consumed all the data. // 2. HybridExtractor will first try to do extraction in log mode, and then choose log or - // tsfile mode to continue extracting, but if (leader data regions num * Wal size) > (maximum - // size of wal buffer), the write operation will be throttled, so we should not extract any - // more tablet events. + // tsfile mode to continue extracting, but if Wal size > maximum size of wal buffer, + // the write operation will be throttled, so we should not extract any more tablet events. // 3. The number of tsfile events in the pending queue has exceeded the limit. event .getTsFileEpoch() @@ -157,6 +154,10 @@ public class PipeRealtimeDataRegionHybridExtractor extends PipeRealtimeDataRegio final TsFileEpoch.State state = event.getTsFileEpoch().getState(this); switch (state) { + case USING_TABLET: + // Though the data in tsfile event has been extracted in tablet mode, we still need to + // extract the tsfile event to help to determine isTsFileEventCountInQueueExceededLimit(). + // The extracted tsfile event will be discarded in supplyTsFileInsertion. case EMPTY: case USING_TSFILE: case USING_BOTH: @@ -177,10 +178,6 @@ public class PipeRealtimeDataRegionHybridExtractor extends PipeRealtimeDataRegio PipeRealtimeDataRegionHybridExtractor.class.getName(), false); } break; - case USING_TABLET: - // All the tablet events have been extracted, so we can ignore the tsFile event. - event.decreaseReferenceCount(PipeRealtimeDataRegionHybridExtractor.class.getName(), false); - break; default: throw new UnsupportedOperationException( String.format( @@ -230,12 +227,6 @@ public class PipeRealtimeDataRegionHybridExtractor extends PipeRealtimeDataRegio > IoTDBDescriptor.getInstance().getConfig().getThrottleThreshold(); } - private boolean tooManyWALPinned() { - return PipeResourceManager.wal().getApproximatePinnedWALCount() - > Math.max(1, PipeAgent.task().getLeaderDataRegionCount()) - * PipeConfig.getInstance().getPipeMaxAllowedPendingTsFileEpochPerDataRegion(); - } - private boolean isTsFileEventCountInQueueExceededLimit() { return pendingQueue.getTsFileInsertionEventCount() + eventCollectorQueueTsFileSize.get() >= PipeConfig.getInstance().getPipeMaxAllowedPendingTsFileEpochPerDataRegion(); diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java index a94c7889a45..89d7610bae5 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java @@ -189,7 +189,7 @@ public class CommonConfig { private boolean pipeAirGapReceiverEnabled = false; private int pipeAirGapReceiverPort = 9780; - private int pipeMaxAllowedPendingTsFileEpochPerDataRegion = 1; + private int pipeMaxAllowedPendingTsFileEpochPerDataRegion = 2; private long pipeMemoryAllocateRetryIntervalMs = 1000; private int pipeMemoryAllocateMaxRetries = 10;
