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 86469d32210 Pipe: Avoid constructing unnecessary pipe subtasks if db 
and pattern dismatch (#13873)
86469d32210 is described below

commit 86469d322108328b7f892955d92ab42e8f7c0fcb
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)
---
 .../db/pipe/agent/task/PipeDataNodeTaskAgent.java  |  8 ++---
 .../agent/task/builder/PipeDataNodeBuilder.java    |  6 ++--
 .../dataregion/DataRegionListeningFilter.java      | 35 +++++++++++++++++++---
 .../datastructure/pattern/IoTDBTreePattern.java    |  9 ++++++
 .../datastructure/pattern/PrefixTreePattern.java   |  9 ++++++
 .../pipe/datastructure/pattern/TreePattern.java    |  7 +++++
 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 7f43701530f..18de48cfc32 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/IoTDBTreePattern.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/IoTDBTreePattern.java
index 6878759a61d..3b8e008594b 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/IoTDBTreePattern.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/IoTDBTreePattern.java
@@ -101,6 +101,15 @@ public class IoTDBTreePattern extends TreePattern {
     }
   }
 
+  @Override
+  public boolean mayOverlapWithDb(final String db) {
+    try {
+      return patternPartialPath.overlapWith(new PartialPath(db + ".**"));
+    } catch (final IllegalPathException e) {
+      return false;
+    }
+  }
+
   @Override
   public boolean mayOverlapWithDevice(final IDeviceID device) {
     try {
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PrefixTreePattern.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PrefixTreePattern.java
index 237e65ffcea..33c584552dd 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PrefixTreePattern.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/PrefixTreePattern.java
@@ -90,6 +90,15 @@ public class PrefixTreePattern extends TreePattern {
     return pattern.length() <= deviceStr.length() && 
deviceStr.startsWith(pattern);
   }
 
+  @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 mayOverlapWithDevice(final IDeviceID device) {
     final String deviceStr = device.toString();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/TreePattern.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/TreePattern.java
index 7358ce32ec5..7ee8f6dd88d 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/TreePattern.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/datastructure/pattern/TreePattern.java
@@ -130,6 +130,13 @@ public abstract class TreePattern {
   /** Check if a device's all measurements are covered by this pattern. */
   public abstract boolean coversDevice(final IDeviceID 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.
    *

Reply via email to