This is an automated email from the ASF dual-hosted git repository.
rong pushed a commit to branch pipe-table-model-3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/pipe-table-model-3 by this
push:
new 4235724e00a Inject model & db message from session & avoid showing
system attrs
4235724e00a is described below
commit 4235724e00a54c08a84d9df61b45b6c6ec756563
Author: Steve Yurong Su <[email protected]>
AuthorDate: Wed Sep 25 19:20:51 2024 +0800
Inject model & db message from session & avoid showing system attrs
---
.../response/pipe/task/PipeTableResp.java | 10 ++++---
.../execution/config/TableConfigTaskVisitor.java | 5 ++++
.../pipe/config/constant/SystemConstant.java | 31 ++++++++++++++++++++++
3 files changed, 43 insertions(+), 3 deletions(-)
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/pipe/task/PipeTableResp.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/pipe/task/PipeTableResp.java
index bb7cf2137d3..69795c145af 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/pipe/task/PipeTableResp.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/pipe/task/PipeTableResp.java
@@ -27,6 +27,7 @@ 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;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMeta;
+import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
import
org.apache.iotdb.confignode.manager.pipe.extractor.ConfigRegionListeningFilter;
import org.apache.iotdb.confignode.rpc.thrift.TGetAllPipeInfoResp;
import org.apache.iotdb.confignode.rpc.thrift.TShowPipeInfo;
@@ -169,9 +170,12 @@ public class PipeTableResp implements DataSet {
staticMeta.getPipeName(),
staticMeta.getCreationTime(),
runtimeMeta.getStatus().get().name(),
- staticMeta.getExtractorParameters().toString(),
- staticMeta.getProcessorParameters().toString(),
- staticMeta.getConnectorParameters().toString(),
+
SystemConstant.copyAndRemoveSystemKeys(staticMeta.getExtractorParameters())
+ .toString(),
+
SystemConstant.copyAndRemoveSystemKeys(staticMeta.getProcessorParameters())
+ .toString(),
+
SystemConstant.copyAndRemoveSystemKeys(staticMeta.getConnectorParameters())
+ .toString(),
exceptionMessageBuilder.toString());
final PipeTemporaryMeta temporaryMeta = pipeMeta.getTemporaryMeta();
final boolean canCalculateOnLocal = canCalculateOnLocal(pipeMeta);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/TableConfigTaskVisitor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/TableConfigTaskVisitor.java
index b39597dc77f..606a0958f0c 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/TableConfigTaskVisitor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/TableConfigTaskVisitor.java
@@ -464,6 +464,11 @@ public class TableConfigTaskVisitor extends
AstVisitor<IConfigTask, MPPQueryCont
protected IConfigTask visitCreatePipe(CreatePipe node, MPPQueryContext
context) {
context.setQueryType(QueryType.WRITE);
+ // Inject table model into the extractor attributes
+ node.getExtractorAttributes()
+ .put(SystemConstant.SESSION_MODEL_KEY,
SystemConstant.SESSION_MODEL_TABLE_VALUE);
+
+ // Inject the database name from the session if it is not specified in the
create pipe statement
final String databaseSetInSession = clientSession.getDatabaseName();
if (Objects.nonNull(databaseSetInSession)) {
final PipeParameters parameters = new
PipeParameters(node.getExtractorAttributes());
diff --git
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/SystemConstant.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/SystemConstant.java
index 923cb7fe342..2e0941b0ea1 100644
---
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/SystemConstant.java
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/SystemConstant.java
@@ -19,13 +19,44 @@
package org.apache.iotdb.commons.pipe.config.constant;
+import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
+
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Set;
+
public class SystemConstant {
public static final String RESTART_KEY = "__system.restart";
public static final boolean RESTART_DEFAULT_VALUE = false;
+ public static final String SESSION_MODEL_KEY = "__system.session.model";
+ public static final String SESSION_MODEL_TREE_VALUE = "tree";
+ public static final String SESSION_MODEL_TABLE_VALUE = "table";
+
public static final String SESSION_DATABASE_KEY =
"__system.session.database";
+ /////////////////////////////////// Utility
///////////////////////////////////
+
+ public static final Set<String> SYSTEM_KEYS = new HashSet<>();
+
+ static {
+ SYSTEM_KEYS.add(RESTART_KEY);
+ SYSTEM_KEYS.add(SESSION_MODEL_KEY);
+ SYSTEM_KEYS.add(SESSION_DATABASE_KEY);
+ }
+
+ public static PipeParameters copyAndRemoveSystemKeys(final PipeParameters
givenPipeParameters) {
+ final Map<String, String> attributes = new
HashMap<>(givenPipeParameters.getAttribute());
+ for (final String key : SYSTEM_KEYS) {
+ attributes.remove(key);
+ }
+ return new PipeParameters(attributes);
+ }
+
+ /////////////////////////////////// Private Constructor
///////////////////////////////////
+
private SystemConstant() {
throw new IllegalStateException("Utility class");
}