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));