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 94062d1a10b [HOTFIX] Pipe: PipeTaskRuntimeConfiguration cannot be cast 
to PipeTaskCollectorRuntimeEnvironment && unreal message when dropping pipe  
(#10121)
94062d1a10b is described below

commit 94062d1a10bddc7a5909ed1572f7eaf8ff3f8dbe
Author: Steve Yurong Su <[email protected]>
AuthorDate: Mon Jun 12 11:15:17 2023 +0800

    [HOTFIX] Pipe: PipeTaskRuntimeConfiguration cannot be cast to 
PipeTaskCollectorRuntimeEnvironment && unreal message when dropping pipe  
(#10121)
---
 .../iotdb/confignode/manager/pipe/task/PipeTaskCoordinator.java      | 5 +++--
 .../org/apache/iotdb/db/pipe/collector/IoTDBDataRegionCollector.java | 3 ++-
 .../org/apache/iotdb/db/pipe/task/stage/PipeTaskCollectorStage.java  | 3 +--
 3 files changed, 6 insertions(+), 5 deletions(-)

diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/task/PipeTaskCoordinator.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/task/PipeTaskCoordinator.java
index 6f47bded042..8719d1abb15 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/task/PipeTaskCoordinator.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/task/PipeTaskCoordinator.java
@@ -73,11 +73,12 @@ public class PipeTaskCoordinator {
   }
 
   public TSStatus dropPipe(String pipeName) {
-    TSStatus status = configManager.getProcedureManager().dropPipe(pipeName);
+    final boolean isPipeExistedBeforeDrop = 
pipeTaskInfo.isPipeExisted(pipeName);
+    final TSStatus status = 
configManager.getProcedureManager().dropPipe(pipeName);
     if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
       LOGGER.warn(String.format("Failed to drop pipe %s. Result status: %s.", 
pipeName, status));
     }
-    return pipeTaskInfo.isPipeExisted(pipeName)
+    return isPipeExistedBeforeDrop
         ? status
         : RpcUtils.getStatus(
             TSStatusCode.PIPE_NOT_EXIST_ERROR,
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/collector/IoTDBDataRegionCollector.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/collector/IoTDBDataRegionCollector.java
index 7916e1d92fb..d1653b30c59 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/collector/IoTDBDataRegionCollector.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/collector/IoTDBDataRegionCollector.java
@@ -137,7 +137,8 @@ public class IoTDBDataRegionCollector implements 
PipeCollector {
   @Override
   public void customize(PipeParameters parameters, 
PipeCollectorRuntimeConfiguration configuration)
       throws Exception {
-    dataRegionId = ((PipeTaskCollectorRuntimeEnvironment) 
configuration).getRegionId();
+    dataRegionId =
+        ((PipeTaskCollectorRuntimeEnvironment) 
configuration.getRuntimeEnvironment()).getRegionId();
 
     historicalCollector.customize(parameters, configuration);
     realtimeCollector.customize(parameters, configuration);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/task/stage/PipeTaskCollectorStage.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/task/stage/PipeTaskCollectorStage.java
index 524b94eb604..6c6f2498b97 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/task/stage/PipeTaskCollectorStage.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/task/stage/PipeTaskCollectorStage.java
@@ -29,7 +29,6 @@ import 
org.apache.iotdb.db.pipe.config.plugin.configuraion.PipeTaskRuntimeConfig
 import 
org.apache.iotdb.db.pipe.config.plugin.env.PipeTaskCollectorRuntimeEnvironment;
 import org.apache.iotdb.db.pipe.task.connection.EventSupplier;
 import org.apache.iotdb.pipe.api.PipeCollector;
-import 
org.apache.iotdb.pipe.api.customizer.configuration.PipeCollectorRuntimeConfiguration;
 import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
 import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
 import org.apache.iotdb.pipe.api.exception.PipeException;
@@ -61,7 +60,7 @@ public class PipeTaskCollectorStage extends PipeTaskStage {
       pipeCollector.validate(new PipeParameterValidator(collectorParameters));
 
       // 2. customize collector
-      final PipeCollectorRuntimeConfiguration runtimeConfiguration =
+      final PipeTaskRuntimeConfiguration runtimeConfiguration =
           new PipeTaskRuntimeConfiguration(
               new PipeTaskCollectorRuntimeEnvironment(
                   pipeName, creationTime, dataRegionId.getId(), pipeTaskMeta));

Reply via email to