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

rong pushed a commit to branch pipe-hybrid-improvement
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit eb5743da994e57da6e0979b42888ce16e8bd2877
Author: Steve Yurong Su <[email protected]>
AuthorDate: Thu Nov 9 13:33:35 2023 +0800

    Pipe: introduce pipe_max_allowed_pinned_memtable_count param for realtime 
hybrid mode
---
 .../PipeRealtimeDataRegionHybridExtractor.java     | 43 ++++++++++++++++------
 .../PipeRealtimeDataRegionLogExtractor.java        |  2 +-
 .../apache/iotdb/commons/conf/CommonConfig.java    |  9 +++++
 .../iotdb/commons/conf/CommonDescriptor.java       |  5 +++
 .../iotdb/commons/pipe/config/PipeConfig.java      |  5 +++
 5 files changed, 51 insertions(+), 13 deletions(-)

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 3ba5a4c496f..91b786bfffa 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,6 +26,7 @@ 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;
@@ -74,16 +75,7 @@ public class PipeRealtimeDataRegionHybridExtractor extends 
PipeRealtimeDataRegio
   }
 
   private void extractTabletInsertion(PipeRealtimeEvent event) {
-    if (!isStartedToSupply
-        || mayWalSizeReachThrottleThreshold()
-        || 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 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.
+    if (canNotUseTabletAnyMore()) {
       event
           .getTsFileEpoch()
           .migrateState(
@@ -223,11 +215,31 @@ public class PipeRealtimeDataRegionHybridExtractor 
extends PipeRealtimeDataRegio
     }
   }
 
+  private boolean canNotUseTabletAnyMore() {
+    // In the following 4 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 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 pinned memtables has reached the dangerous threshold.
+    //  4. The number of tsfile events in the pending queue has exceeded the 
limit.
+    return !isStartedToSupply
+        || mayWalSizeReachThrottleThreshold()
+        || mayMemTablePinnedCountReachDangerousThreshold()
+        || isTsFileEventCountInQueueExceededLimit();
+  }
+
   private boolean mayWalSizeReachThrottleThreshold() {
     return 3 * WALManager.getInstance().getTotalDiskUsage()
         > IoTDBDescriptor.getInstance().getConfig().getThrottleThreshold();
   }
 
+  private boolean mayMemTablePinnedCountReachDangerousThreshold() {
+    return PipeResourceManager.wal().getPinnedWalCount()
+        >= PipeConfig.getInstance().getPipeMaxAllowedPinnedMemTableCount();
+  }
+
   private boolean isTsFileEventCountInQueueExceededLimit() {
     return pendingQueue.getTsFileInsertionEventCount()
             + processorEventCollectorQueueTsFileSize.get()
@@ -286,8 +298,15 @@ public class PipeRealtimeDataRegionHybridExtractor extends 
PipeRealtimeDataRegio
         .getTsFileEpoch()
         .migrateState(
             this,
-            state ->
-                (state.equals(TsFileEpoch.State.EMPTY)) ? 
TsFileEpoch.State.USING_TABLET : state);
+            state -> {
+              if (!state.equals(TsFileEpoch.State.EMPTY)) {
+                return state;
+              }
+
+              return canNotUseTabletAnyMore()
+                  ? TsFileEpoch.State.USING_TSFILE
+                  : TsFileEpoch.State.USING_TABLET;
+            });
 
     final TsFileEpoch.State state = event.getTsFileEpoch().getState(this);
     switch (state) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionLogExtractor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionLogExtractor.java
index c309bd2c620..2d857bcad9b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionLogExtractor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionLogExtractor.java
@@ -50,7 +50,7 @@ public class PipeRealtimeDataRegionLogExtractor extends 
PipeRealtimeDataRegionEx
     } else {
       throw new UnsupportedOperationException(
           String.format(
-              "Unsupported event type %s for hybrid realtime extractor %s",
+              "Unsupported event type %s for log realtime extractor %s",
               eventToExtract.getClass(), this));
     }
   }
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 48a3a30dda9..02b81938c7d 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
@@ -184,6 +184,7 @@ public class CommonConfig {
   private int pipeAirGapReceiverPort = 9780;
 
   private int pipeMaxAllowedPendingTsFileEpochPerDataRegion = 2;
+  private int pipeMaxAllowedPinnedMemTableCount = 50;
 
   private boolean pipeMemoryManagementEnabled = true;
   private long pipeMemoryAllocateRetryIntervalMs = 1000;
@@ -741,6 +742,14 @@ public class CommonConfig {
     this.pipeMaxAllowedPendingTsFileEpochPerDataRegion = 
pipeExtractorPendingQueueTsfileLimit;
   }
 
+  public int getPipeMaxAllowedPinnedMemTableCount() {
+    return pipeMaxAllowedPinnedMemTableCount;
+  }
+
+  public void setPipeMaxAllowedPinnedMemTableCount(int 
pipeMaxAllowedPinnedMemTableCount) {
+    this.pipeMaxAllowedPinnedMemTableCount = pipeMaxAllowedPinnedMemTableCount;
+  }
+
   public boolean getPipeMemoryManagementEnabled() {
     return pipeMemoryManagementEnabled;
   }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
index cc794954a78..fc35c1e9fcc 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
@@ -392,6 +392,11 @@ public class CommonDescriptor {
             properties.getProperty(
                 "pipe_max_allowed_pending_tsfile_epoch_per_data_region",
                 
String.valueOf(config.getPipeMaxAllowedPendingTsFileEpochPerDataRegion()))));
+    config.setPipeMaxAllowedPinnedMemTableCount(
+        Integer.parseInt(
+            properties.getProperty(
+                "pipe_max_allowed_pinned_memtable_count",
+                
String.valueOf(config.getPipeMaxAllowedPinnedMemTableCount()))));
 
     config.setPipeMemoryManagementEnabled(
         Boolean.parseBoolean(
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
index 47e3f1d5003..f2f018f703e 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
@@ -161,6 +161,10 @@ public class PipeConfig {
     return COMMON_CONFIG.getPipeMaxAllowedPendingTsFileEpochPerDataRegion();
   }
 
+  public int getPipeMaxAllowedPinnedMemTableCount() {
+    return COMMON_CONFIG.getPipeMaxAllowedPinnedMemTableCount();
+  }
+
   /////////////////////////////// Memory ///////////////////////////////
 
   public boolean getPipeMemoryManagementEnabled() {
@@ -244,6 +248,7 @@ public class PipeConfig {
     LOGGER.info(
         "PipeMaxAllowedPendingTsFileEpochPerDataRegion: {}",
         getPipeMaxAllowedPendingTsFileEpochPerDataRegion());
+    LOGGER.info("PipeMaxAllowedPinnedMemTableCount: {}", 
getPipeMaxAllowedPinnedMemTableCount());
 
     LOGGER.info("PipeMemoryManagementEnabled: {}", 
getPipeMemoryManagementEnabled());
     LOGGER.info("PipeMemoryAllocateMaxRetries: {}", 
getPipeMemoryAllocateMaxRetries());

Reply via email to