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.