This is an automated email from the ASF dual-hosted git repository.

rong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 06a44df157e Pipe: fine tune hybrid mode by removing tooManyWALPinned 
judgement and increasing pipeMaxAllowedPendingTsFileEpochPerDataRegion to 2 
(#11382)
06a44df157e is described below

commit 06a44df157ec946e812ec58d464e1306bb253aa5
Author: Steve Yurong Su <[email protected]>
AuthorDate: Wed Oct 25 12:53:40 2023 +0800

    Pipe: fine tune hybrid mode by removing tooManyWALPinned judgement and 
increasing pipeMaxAllowedPendingTsFileEpochPerDataRegion to 2 (#11382)
---
 .../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;

Reply via email to