diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java index 88cd54dbe28c..0710fe345b85 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java @@ -908,7 +908,10 @@ public static PipeMeta copyAndFilterOutNonWorkingDataRegionPipeTasks(PipeMeta or try { return !DataRegionListeningFilter.shouldDatabaseBeListened( - copiedPipeMeta.getStaticMeta().getSourceParameters(), isTableModel, database); + copiedPipeMeta.getStaticMeta().getSourceParameters(), + isTableModel, + database, + copiedPipeMeta.getStaticMeta().getPipeType()); } catch (final Exception e) { return false; } 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 c47c54f675a7..5241ea0c916d 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 @@ -161,7 +161,7 @@ protected void createPipeTask( final boolean needConstructDataRegionTask = StorageEngine.getInstance().getAllDataRegionIds().contains(dataRegionId) && DataRegionListeningFilter.shouldDataRegionBeListened( - sourceParameters, dataRegionId); + sourceParameters, dataRegionId, pipeStaticMeta.getPipeType()); final boolean needConstructSchemaRegionTask = SchemaEngine.getInstance() .getAllSchemaRegionIds() @@ -861,12 +861,21 @@ public ProgressIndex getPipeTaskProgressIndex(final String pipeName, final int c throw new PipeException(DataNodePipeMessages.PIPE_META_NOT_FOUND + pipeName); } - return pipeMetaKeeper - .getPipeMeta(pipeName) - .getRuntimeMeta() - .getConsensusGroupId2TaskMetaMap() - .get(consensusGroupId) - .getProgressIndex(); + final PipeTaskMeta pipeTaskMeta = + pipeMetaKeeper + .getPipeMeta(pipeName) + .getRuntimeMeta() + .getConsensusGroupId2TaskMetaMap() + .get(consensusGroupId); + if (pipeTaskMeta == null) { + throw new PipeException( + String.format( + DataNodePipeMessages + .PIPE_EXCEPTION_FAILED_TO_GET_PIPE_TASK_PROGRESS_INDEX_WITH_PIPE_NAME_S_CFE9DE7C, + pipeName, + consensusGroupId)); + } + return pipeTaskMeta.getProgressIndex(); } finally { releaseReadLock(); } @@ -951,7 +960,7 @@ private List> collectPipeTasksToBeCreated( final boolean needConstructDataRegionTask = dataRegionIds.contains(dataRegionId) && DataRegionListeningFilter.shouldDataRegionBeListened( - sourceParameters, dataRegionId); + sourceParameters, dataRegionId, pipeStaticMeta.getPipeType()); final boolean needConstructSchemaRegionTask = schemaRegionIds.contains(new SchemaRegionId(consensusGroupId)) && SchemaRegionListeningFilter.shouldSchemaRegionBeListened( 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 46a10135d886..2896ff6039a5 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 @@ -68,7 +68,7 @@ public Map buildTasksWithInternalSource() throws IllegalPathE final boolean needConstructDataRegionTask = dataRegionIds.contains(dataRegionId) && DataRegionListeningFilter.shouldDataRegionBeListened( - sourceParameters, dataRegionId); + sourceParameters, dataRegionId, pipeStaticMeta.getPipeType()); final boolean needConstructSchemaRegionTask = schemaRegionIds.contains(new SchemaRegionId(consensusGroupId)) && SchemaRegionListeningFilter.shouldSchemaRegionBeListened( diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilter.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilter.java index ad4c1fddff3e..bbaf9ff97457 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilter.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilter.java @@ -23,6 +23,7 @@ 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.agent.task.meta.PipeType; 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; @@ -58,9 +59,12 @@ public class DataRegionListeningFilter { } public static boolean shouldDatabaseBeListened( - final PipeParameters parameters, final boolean isTableModel, final String databaseRawName) + final PipeParameters parameters, + final boolean isTableModel, + final String databaseRawName, + final PipeType pipeType) throws IllegalPathException { - if (isAuditDatabase(databaseRawName)) { + if (!PipeType.CONSENSUS.equals(pipeType) && isAuditDatabase(databaseRawName)) { return false; } @@ -90,7 +94,8 @@ public static boolean shouldDatabaseBeListened( } public static boolean shouldDataRegionBeListened( - PipeParameters parameters, DataRegionId dataRegionId) throws IllegalPathException { + final PipeParameters parameters, final DataRegionId dataRegionId, final PipeType pipeType) + throws IllegalPathException { final Pair insertionDeletionListeningOptionPair = parseInsertionDeletionListeningOptionPair(parameters); final boolean hasSpecificListeningOption = @@ -106,7 +111,7 @@ public static boolean shouldDataRegionBeListened( } final String databaseRawName = dataRegion.getDatabaseName(); - if (isAuditDatabase(databaseRawName)) { + if (!PipeType.CONSENSUS.equals(pipeType) && isAuditDatabase(databaseRawName)) { return false; } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgentTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgentTest.java index 3dd93e877890..018ec371d91d 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgentTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgentTest.java @@ -23,7 +23,9 @@ import org.apache.iotdb.commons.consensus.index.ProgressIndex; import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex; import org.apache.iotdb.commons.consensus.index.impl.SimpleProgressIndex; +import org.apache.iotdb.commons.pipe.agent.task.PipeTaskAgent; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMetaKeeper; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeRuntimeMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; @@ -34,6 +36,7 @@ import org.junit.Assert; import org.junit.Test; +import java.lang.reflect.Field; import java.util.HashMap; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -44,6 +47,26 @@ public class PipeDataNodeTaskAgentTest { private static final int LOCAL_NODE_ID = 1; private static final int REGION_ID = 7; + @Test + public void testGetPipeTaskProgressIndexReportsMissingTaskMeta() throws Exception { + final PipeDataNodeTaskAgent taskAgent = new PipeDataNodeTaskAgent(); + final Field pipeMetaKeeperField = PipeTaskAgent.class.getDeclaredField("pipeMetaKeeper"); + pipeMetaKeeperField.setAccessible(true); + final PipeMetaKeeper pipeMetaKeeper = (PipeMetaKeeper) pipeMetaKeeperField.get(taskAgent); + + final String pipeName = PipeStaticMeta.CONSENSUS_PIPE_PREFIX + "DataRegion[7]_1_2"; + pipeMetaKeeper.addPipeMeta( + new PipeMeta( + new PipeStaticMeta(pipeName, 1L, new HashMap<>(), new HashMap<>(), new HashMap<>()), + new PipeRuntimeMeta())); + + final PipeException exception = + Assert.assertThrows( + PipeException.class, () -> taskAgent.getPipeTaskProgressIndex(pipeName, REGION_ID)); + Assert.assertTrue(exception.getMessage().contains(pipeName)); + Assert.assertTrue(exception.getMessage().contains(String.valueOf(REGION_ID))); + } + @Test public void testCreateMemoryCheckStillRunsWhenNoPipeTasksNeedToBeCreated() throws Exception { final boolean originalPipeEnableMemoryCheck = diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilterTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilterTest.java index 4901afcf9120..4257b032ef42 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilterTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilterTest.java @@ -19,6 +19,7 @@ package org.apache.iotdb.db.pipe.source.dataregion; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeType; import org.apache.iotdb.commons.pipe.config.constant.SystemConstant; import org.apache.iotdb.commons.subscription.meta.topic.TopicMeta; import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters; @@ -34,17 +35,31 @@ public class DataRegionListeningFilterTest { @Test - public void testAuditDatabaseIsNeverListened() throws Exception { + public void testAuditDatabaseIsOnlyListenedByConsensusPipes() throws Exception { final Map topicAttributes = new HashMap<>(); topicAttributes.put(SystemConstant.SQL_DIALECT_KEY, SystemConstant.SQL_DIALECT_TABLE_VALUE); final PipeParameters parameters = new PipeParameters( new TopicMeta("topic", 1, topicAttributes).generateExtractorAttributes("root")); - assertFalse(DataRegionListeningFilter.shouldDatabaseBeListened(parameters, true, "__audit")); assertFalse( - DataRegionListeningFilter.shouldDatabaseBeListened(parameters, false, "root.__audit")); - assertFalse(DataRegionListeningFilter.shouldDatabaseBeListened(parameters, true, "__AUDIT")); - assertTrue(DataRegionListeningFilter.shouldDatabaseBeListened(parameters, true, "user_db")); + DataRegionListeningFilter.shouldDatabaseBeListened( + parameters, true, "__audit", PipeType.SUBSCRIPTION)); + assertFalse( + DataRegionListeningFilter.shouldDatabaseBeListened( + parameters, false, "root.__audit", PipeType.SUBSCRIPTION)); + assertFalse( + DataRegionListeningFilter.shouldDatabaseBeListened( + parameters, true, "__AUDIT", PipeType.SUBSCRIPTION)); + assertTrue( + DataRegionListeningFilter.shouldDatabaseBeListened( + parameters, true, "user_db", PipeType.SUBSCRIPTION)); + + assertTrue( + DataRegionListeningFilter.shouldDatabaseBeListened( + parameters, true, "__audit", PipeType.CONSENSUS)); + assertFalse( + DataRegionListeningFilter.shouldDatabaseBeListened( + parameters, true, "__audit", PipeType.USER)); } }