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