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.
*