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

Reply via email to