This is an automated email from the ASF dual-hosted git repository.

rong pushed a commit to branch select-into
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/select-into by this push:
     new da0da6c  SelectIntoTransporter
da0da6c is described below

commit da0da6c5fe002f5c5114583bb73e919c28309066
Author: Steve Yurong Su <[email protected]>
AuthorDate: Thu Jul 15 20:43:33 2021 +0800

    SelectIntoTransporter
---
 ...sporter.java => InsertTabletPlanGenerator.java} | 36 ++++++----
 .../engine/transporter/SelectIntoTransporter.java  | 84 +++++++++++++++++++++-
 .../db/qp/logical/crud/SelectIntoOperator.java     |  2 +
 3 files changed, 108 insertions(+), 14 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/transporter/SelectIntoTransporter.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/transporter/InsertTabletPlanGenerator.java
similarity index 58%
copy from 
server/src/main/java/org/apache/iotdb/db/engine/transporter/SelectIntoTransporter.java
copy to 
server/src/main/java/org/apache/iotdb/db/engine/transporter/InsertTabletPlanGenerator.java
index d9749d0..c631b10 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/transporter/SelectIntoTransporter.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/transporter/InsertTabletPlanGenerator.java
@@ -19,31 +19,43 @@
 
 package org.apache.iotdb.db.engine.transporter;
 
-import org.apache.iotdb.db.exception.IoTDBException;
 import org.apache.iotdb.db.metadata.PartialPath;
+import org.apache.iotdb.db.qp.physical.crud.InsertTabletPlan;
+import org.apache.iotdb.tsfile.read.common.RowRecord;
 import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet;
+import org.apache.iotdb.tsfile.write.record.Tablet;
 
-import java.io.IOException;
+import java.util.ArrayList;
 import java.util.List;
 
-public class SelectIntoTransporter {
+public class InsertTabletPlanGenerator {
 
   private final QueryDataSet queryDataSet;
   private final List<PartialPath> intoPaths;
-  private final int fetchSize;
 
-  public SelectIntoTransporter(
-      QueryDataSet queryDataSet, List<PartialPath> intoPaths, int fetchSize) {
+  private final String device;
+  private final List<Integer> measurementIdIndexes;
+
+  private Tablet tablet;
+
+  public InsertTabletPlanGenerator(
+      QueryDataSet queryDataSet, List<PartialPath> intoPaths, String device) {
     this.queryDataSet = queryDataSet;
     this.intoPaths = intoPaths;
-    this.fetchSize = fetchSize;
+
+    this.device = device;
+    measurementIdIndexes = new ArrayList<>();
   }
 
-  public void transport() throws IoTDBException, IOException {
-    while (queryDataSet.hasNext()) {
-      transportTablets();
-    }
+  public void addMeasurementIdIndex(int measurementIdIndex) {
+    measurementIdIndexes.add(measurementIdIndex);
   }
 
-  private void transportTablets() {}
+  public void constructNewTablet() {}
+
+  public void collectRowRecord(RowRecord rowRecord) {}
+
+  public InsertTabletPlan getInsertTabletPlan() {
+    throw new UnsupportedOperationException();
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/transporter/SelectIntoTransporter.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/transporter/SelectIntoTransporter.java
index d9749d0..f0d4718 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/transporter/SelectIntoTransporter.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/transporter/SelectIntoTransporter.java
@@ -20,18 +20,31 @@
 package org.apache.iotdb.db.engine.transporter;
 
 import org.apache.iotdb.db.exception.IoTDBException;
+import org.apache.iotdb.db.exception.metadata.IllegalPathException;
 import org.apache.iotdb.db.metadata.PartialPath;
+import org.apache.iotdb.db.qp.physical.crud.InsertMultiTabletPlan;
+import org.apache.iotdb.db.qp.physical.crud.InsertTabletPlan;
+import org.apache.iotdb.tsfile.read.common.RowRecord;
 import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet;
 
 import java.io.IOException;
+import java.util.ArrayList;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
 
 public class SelectIntoTransporter {
 
+  private static final Pattern leveledPathNodePattern = 
Pattern.compile("\\$\\{\\w+}");
+
   private final QueryDataSet queryDataSet;
   private final List<PartialPath> intoPaths;
   private final int fetchSize;
 
+  private InsertTabletPlanGenerator[] insertTabletPlanGenerators;
+
   public SelectIntoTransporter(
       QueryDataSet queryDataSet, List<PartialPath> intoPaths, int fetchSize) {
     this.queryDataSet = queryDataSet;
@@ -40,10 +53,77 @@ public class SelectIntoTransporter {
   }
 
   public void transport() throws IoTDBException, IOException {
+    generateActualIntoPaths();
+    constructTabletGenerators();
+    doTransport();
+  }
+
+  private void generateActualIntoPaths() throws IllegalPathException {
+    for (int i = 0; i < intoPaths.size(); ++i) {
+      intoPaths.set(i, generateActualIntoPath(i));
+    }
+  }
+
+  private PartialPath generateActualIntoPath(int index) throws 
IllegalPathException {
+    String[] nodes = intoPaths.get(index).getNodes();
+
+    int indexOfLeftBracket = nodes[0].indexOf("(");
+    if (indexOfLeftBracket != -1) {
+      nodes[0] = nodes[0].substring(indexOfLeftBracket + 1);
+    }
+    int indexOfRightBracket = nodes[nodes.length - 1].indexOf(")");
+    if (indexOfRightBracket != -1) {
+      nodes[nodes.length - 1] = nodes[nodes.length - 1].substring(0, 
indexOfRightBracket);
+    }
+
+    StringBuffer sb = new StringBuffer();
+    Matcher m = 
leveledPathNodePattern.matcher(queryDataSet.getPaths().get(index).getFullPath());
+    while (m.find()) {
+      String param = m.group();
+      String value = nodes[Integer.parseInt(param.substring(2, param.length() 
- 1).trim())];
+      m.appendReplacement(sb, value == null ? "" : value);
+    }
+    m.appendTail(sb);
+    return new PartialPath(sb.toString());
+  }
+
+  private void constructTabletGenerators() {
+    Map<String, InsertTabletPlanGenerator> deviceTabletGeneratorMap = new 
HashMap<>();
+    for (int i = 0, intoPathsSize = intoPaths.size(); i < intoPathsSize; i++) {
+      String device = intoPaths.get(i).getDevice();
+      if (!deviceTabletGeneratorMap.containsKey(device)) {
+        deviceTabletGeneratorMap.put(
+            device, new InsertTabletPlanGenerator(queryDataSet, intoPaths, 
device));
+      }
+      deviceTabletGeneratorMap.get(device).addMeasurementIdIndex(i);
+    }
+    insertTabletPlanGenerators =
+        deviceTabletGeneratorMap.values().toArray(new 
InsertTabletPlanGenerator[0]);
+  }
+
+  private void doTransport() throws IOException {
     while (queryDataSet.hasNext()) {
-      transportTablets();
+      List<InsertTabletPlan> insertTabletPlanList = new ArrayList<>();
+      for (InsertTabletPlanGenerator insertTabletPlanGenerator : 
insertTabletPlanGenerators) {
+        insertTabletPlanGenerator.constructNewTablet();
+      }
+      collectRowRecordIntoInsertTabletPlanGenerators();
+      for (InsertTabletPlanGenerator insertTabletPlanGenerator : 
insertTabletPlanGenerators) {
+        
insertTabletPlanList.add(insertTabletPlanGenerator.getInsertTabletPlan());
+      }
+      InsertMultiTabletPlan insertMultiTabletPlan = new 
InsertMultiTabletPlan(insertTabletPlanList);
+      // TODO: execute insertMultiTabletPlan
     }
   }
 
-  private void transportTablets() {}
+  private void collectRowRecordIntoInsertTabletPlanGenerators() throws 
IOException {
+    int count = 0;
+    while (queryDataSet.hasNext() && count < fetchSize) {
+      RowRecord rowRecord = queryDataSet.next();
+      for (InsertTabletPlanGenerator insertTabletPlanGenerator : 
insertTabletPlanGenerators) {
+        insertTabletPlanGenerator.collectRowRecord(rowRecord);
+      }
+      ++count;
+    }
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/qp/logical/crud/SelectIntoOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/qp/logical/crud/SelectIntoOperator.java
index ee9f7fb..1859cef 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/qp/logical/crud/SelectIntoOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/qp/logical/crud/SelectIntoOperator.java
@@ -50,6 +50,8 @@ public class SelectIntoOperator extends Operator {
 
   public void check() throws LogicalOperatorException {
     queryOperator.check();
+
+    // TODO: check query plan type
   }
 
   public void setQueryOperator(QueryOperator queryOperator) {

Reply via email to