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

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


The following commit(s) were added to refs/heads/rc/1.3.3 by this push:
     new d937c997879 Pipe: Avoid constructing unnecessary pipe subtasks if db 
and pattern dismatch (#13873) (#13892)
d937c997879 is described below

commit d937c9978793c608e5a61ed0709f6fdc7d830779
Author: Steve Yurong Su <[email protected]>
AuthorDate: Thu Oct 24 11:36:27 2024 +0800

    Pipe: Avoid constructing unnecessary pipe subtasks if db and pattern 
dismatch (#13873) (#13892)
    
    * Pipe: Avoid constructing unnecessary pipe subtasks if db and pattern 
dismatch (#13873)
    
    (cherry picked from commit 86469d322108328b7f892955d92ab42e8f7c0fcb)
    
    * Update DataRegionListeningFilter.java
---
 .../db/pipe/agent/task/PipeDataNodeTaskAgent.java  |  8 ++++----
 .../agent/task/builder/PipeDataNodeBuilder.java    |  6 ++++--
 .../dataregion/DataRegionListeningFilter.java      | 23 ++++++++++++++++++----
 .../datastructure/pattern/IoTDBPipePattern.java    |  9 +++++++++
 .../pipe/datastructure/pattern/PipePattern.java    |  7 +++++++
 .../datastructure/pattern/PrefixPipePattern.java   |  9 +++++++++
 6 files changed, 52 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..c78d844f0b1 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,12 @@
 
 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.PipePattern;
+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 +60,24 @@ 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;
+    }
+
+    return PipePattern.parsePipePatternFromSourceParameters(parameters)
+        .mayOverlapWithDb(dataRegion.getDatabaseName());
   }
 
   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