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

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

commit 851cfbc8de2d1b584f0f4f97afcc3f52332aa453
Author: Steve Yurong Su <[email protected]>
AuthorDate: Wed Oct 23 18:45:45 2024 +0800

    Pipe: Avoid constructing unnecessary pipe subtasks if db and pattern 
dismatch (#13873)
    
    (cherry picked from commit 86469d322108328b7f892955d92ab42e8f7c0fcb)
---
 .../db/pipe/agent/task/PipeDataNodeTaskAgent.java  |  8 ++---
 .../agent/task/builder/PipeDataNodeBuilder.java    |  6 ++--
 .../dataregion/DataRegionListeningFilter.java      | 35 +++++++++++++++++++---
 .../datastructure/pattern/IoTDBPipePattern.java    |  9 ++++++
 .../pipe/datastructure/pattern/PipePattern.java    |  7 +++++
 .../datastructure/pattern/PrefixPipePattern.java   |  9 ++++++
 6 files changed, 64 insertions(+), 10 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
index 4ef91b9e37a..f55da4be67b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
@@ -118,11 +118,11 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
       throws IllegalPathException {
     if (pipeTaskMeta.getLeaderNodeId() == CONFIG.getDataNodeId()) {
       final PipeParameters extractorParameters = 
pipeStaticMeta.getExtractorParameters();
+      final DataRegionId dataRegionId = new DataRegionId(consensusGroupId);
       final boolean needConstructDataRegionTask =
-          StorageEngine.getInstance()
-                  .getAllDataRegionIds()
-                  .contains(new DataRegionId(consensusGroupId))
-              && 
DataRegionListeningFilter.shouldDataRegionBeListened(extractorParameters);
+          
StorageEngine.getInstance().getAllDataRegionIds().contains(dataRegionId)
+              && DataRegionListeningFilter.shouldDataRegionBeListened(
+                  extractorParameters, dataRegionId);
       final boolean needConstructSchemaRegionTask =
           SchemaEngine.getInstance()
                   .getAllSchemaRegionIds()
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeBuilder.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeBuilder.java
index 49c55f2dc96..3d639215258 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeBuilder.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeBuilder.java
@@ -64,9 +64,11 @@ public class PipeDataNodeBuilder {
 
       if (pipeTaskMeta.getLeaderNodeId() == CONFIG.getDataNodeId()) {
         final PipeParameters extractorParameters = 
pipeStaticMeta.getExtractorParameters();
+        final DataRegionId dataRegionId = new DataRegionId(consensusGroupId);
         final boolean needConstructDataRegionTask =
-            dataRegionIds.contains(new DataRegionId(consensusGroupId))
-                && 
DataRegionListeningFilter.shouldDataRegionBeListened(extractorParameters);
+            dataRegionIds.contains(dataRegionId)
+                && DataRegionListeningFilter.shouldDataRegionBeListened(
+                    extractorParameters, dataRegionId);
         final boolean needConstructSchemaRegionTask =
             schemaRegionIds.contains(new SchemaRegionId(consensusGroupId))
                 && 
!SchemaRegionListeningFilter.parseListeningPlanTypeSet(extractorParameters)
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/DataRegionListeningFilter.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/DataRegionListeningFilter.java
index a54f5156048..82893ce2a1b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/DataRegionListeningFilter.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/DataRegionListeningFilter.java
@@ -19,9 +19,13 @@
 
 package org.apache.iotdb.db.pipe.extractor.dataregion;
 
+import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.commons.exception.IllegalPathException;
 import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.commons.pipe.agent.task.PipeTask;
+import org.apache.iotdb.commons.pipe.datastructure.pattern.TablePattern;
+import org.apache.iotdb.commons.pipe.datastructure.pattern.TreePattern;
+import org.apache.iotdb.db.storageengine.StorageEngine;
 import org.apache.iotdb.db.storageengine.dataregion.DataRegion;
 import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
 
@@ -57,12 +61,35 @@ public class DataRegionListeningFilter {
     }
   }
 
-  public static boolean shouldDataRegionBeListened(PipeParameters parameters)
-      throws IllegalPathException {
+  public static boolean shouldDataRegionBeListened(
+      PipeParameters parameters, DataRegionId dataRegionId) throws 
IllegalPathException {
     final Pair<Boolean, Boolean> insertionDeletionListeningOptionPair =
         parseInsertionDeletionListeningOptionPair(parameters);
-    return insertionDeletionListeningOptionPair.getLeft()
-        || insertionDeletionListeningOptionPair.getRight();
+    final boolean hasSpecificListeningOption =
+        insertionDeletionListeningOptionPair.getLeft()
+            || insertionDeletionListeningOptionPair.getRight();
+    if (!hasSpecificListeningOption) {
+      return false;
+    }
+
+    final DataRegion dataRegion = 
StorageEngine.getInstance().getDataRegion(dataRegionId);
+    if (dataRegion == null) {
+      return true;
+    }
+
+    final String databaseRawName = dataRegion.getDatabaseName();
+    final String databaseTreeModel =
+        databaseRawName.startsWith("root.") ? databaseRawName : "root." + 
databaseRawName;
+    final String databaseTableModel =
+        databaseRawName.startsWith("root.") ? databaseRawName.substring(5) : 
databaseRawName;
+
+    final TreePattern treePattern = 
TreePattern.parsePipePatternFromSourceParameters(parameters);
+    final TablePattern tablePattern = 
TablePattern.parsePipePatternFromSourceParameters(parameters);
+
+    return treePattern.isTreeModelDataAllowedToBeCaptured()
+            && treePattern.mayOverlapWithDb(databaseTreeModel)
+        || tablePattern.isTableModelDataAllowedToBeCaptured()
+            && tablePattern.matchesDatabase(databaseTableModel);
   }
 
   public static Pair<Boolean, Boolean> 
parseInsertionDeletionListeningOptionPair(
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/IoTDBPipePattern.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/IoTDBPipePattern.java
index 5db5d113c01..ed43952d3ad 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/IoTDBPipePattern.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/IoTDBPipePattern.java
@@ -105,6 +105,15 @@ public class IoTDBPipePattern extends PipePattern {
     }
   }
 
+  @Override
+  public boolean mayOverlapWithDb(final String db) {
+    try {
+      return patternPartialPath.overlapWith(new PartialPath(db + ".**"));
+    } catch (final IllegalPathException e) {
+      return false;
+    }
+  }
+
   @Override
   public boolean matchesMeasurement(final String device, final String 
measurement) {
     // For aligned timeseries, empty measurement is an alias of the time 
column.
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PipePattern.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PipePattern.java
index b6d37fa723e..66e476dabac 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PipePattern.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PipePattern.java
@@ -110,6 +110,13 @@ public abstract class PipePattern {
   /** Check if a device's all measurements are covered by this pattern. */
   public abstract boolean coversDevice(final String device);
 
+  /**
+   * Check if a database may have some measurements matched by the pattern.
+   *
+   * @return {@code true} if the pattern may overlap with the database, {@code 
false} otherwise.
+   */
+  public abstract boolean mayOverlapWithDb(final String db);
+
   /**
    * Check if a device may have some measurements matched by the pattern.
    *
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PrefixPipePattern.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PrefixPipePattern.java
index 33bdf471a0e..349e01ac903 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PrefixPipePattern.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PrefixPipePattern.java
@@ -96,6 +96,15 @@ public class PrefixPipePattern extends PipePattern {
         || (pattern.length() > device.length() && pattern.startsWith(device));
   }
 
+  @Override
+  public boolean mayOverlapWithDb(final String db) {
+    return
+    // for example, pattern is root.a.b and db is root.a.b.c
+    (pattern.length() <= db.length() && db.startsWith(pattern))
+        // for example, pattern is root.a.b.c and db is root.a.b
+        || (pattern.length() > db.length() && pattern.startsWith(db));
+  }
+
   @Override
   public boolean matchesMeasurement(final String device, final String 
measurement) {
     // We assume that the device is already matched.

Reply via email to