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.
