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

haonan pushed a commit to branch optiNewStandaloneAutoCreateSchema
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit b049eba2722f3901697e785fdbe192f2faf9643e
Author: HTHou <[email protected]>
AuthorDate: Thu Jun 9 10:28:36 2022 +0800

    Optimize new standalone auto create schema
---
 .../mpp/plan/analyze/StandaloneSchemaFetcher.java  | 100 +++++++--------------
 .../db/mpp/plan/StandaloneCoordinatorTest.java     |   1 +
 .../plan/scheduler/StandaloneSchedulerTest.java    |  58 +++++++++---
 3 files changed, 80 insertions(+), 79 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/StandaloneSchemaFetcher.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/StandaloneSchemaFetcher.java
index efff9ea7e3..a4fe1616fa 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/StandaloneSchemaFetcher.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/StandaloneSchemaFetcher.java
@@ -18,7 +18,6 @@
  */
 package org.apache.iotdb.db.mpp.plan.analyze;
 
-import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.consensus.SchemaRegionId;
 import org.apache.iotdb.commons.exception.MetadataException;
 import org.apache.iotdb.commons.partition.SchemaPartition;
@@ -32,12 +31,8 @@ import 
org.apache.iotdb.db.mpp.common.schematree.DeviceSchemaInfo;
 import org.apache.iotdb.db.mpp.common.schematree.PathPatternTree;
 import org.apache.iotdb.db.mpp.common.schematree.SchemaTree;
 import org.apache.iotdb.db.mpp.plan.Coordinator;
-import org.apache.iotdb.db.mpp.plan.execution.ExecutionResult;
-import org.apache.iotdb.db.mpp.plan.statement.Statement;
-import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateAlignedTimeSeriesStatement;
-import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateTimeSeriesByDeviceStatement;
-import org.apache.iotdb.db.query.control.SessionManager;
-import org.apache.iotdb.rpc.TSStatusCode;
+import org.apache.iotdb.db.qp.physical.sys.CreateAlignedTimeSeriesPlan;
+import org.apache.iotdb.db.qp.physical.sys.CreateTimeSeriesPlan;
 import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
 import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
@@ -253,68 +248,41 @@ public class StandaloneSchemaFetcher implements 
ISchemaFetcher {
       List<String> measurements,
       List<TSDataType> tsDataTypes,
       boolean isAligned) {
-
-    if (isAligned) {
-      CreateAlignedTimeSeriesStatement createAlignedTimeSeriesStatement =
-          new CreateAlignedTimeSeriesStatement();
-      createAlignedTimeSeriesStatement.setDevicePath(devicePath);
-      createAlignedTimeSeriesStatement.setMeasurements(measurements);
-      createAlignedTimeSeriesStatement.setDataTypes(tsDataTypes);
-      List<TSEncoding> encodings = new ArrayList<>();
-      List<CompressionType> compressors = new ArrayList<>();
-      for (TSDataType dataType : tsDataTypes) {
-        encodings.add(getDefaultEncoding(dataType));
-        
compressors.add(TSFileDescriptor.getInstance().getConfig().getCompressor());
-      }
-      createAlignedTimeSeriesStatement.setEncodings(encodings);
-      createAlignedTimeSeriesStatement.setCompressors(compressors);
-      createAlignedTimeSeriesStatement.setAliasList(null);
-
-      executeCreateStatement(createAlignedTimeSeriesStatement);
-    } else {
-
-      executeCreateTimeseriesByDeviceStatement(
-          new CreateTimeSeriesByDeviceStatement(devicePath, measurements, 
tsDataTypes));
-    }
-  }
-
-  private void executeCreateStatement(Statement statement) {
-    long queryId = SessionManager.getInstance().requestQueryId(false);
-    ExecutionResult executionResult =
-        coordinator.execute(statement, queryId, null, "", partitionFetcher, 
this);
-    // TODO: throw exception
-    try {
-      int statusCode = executionResult.status.getCode();
-      if (statusCode != TSStatusCode.SUCCESS_STATUS.getStatusCode()
-          && statusCode != 
TSStatusCode.PATH_ALREADY_EXIST_ERROR.getStatusCode()) {
-        throw new RuntimeException(
-            "cannot auto create schema, status is: " + executionResult.status);
-      }
-    } finally {
-      coordinator.getQueryExecution(queryId).stopAndCleanup();
-    }
-  }
-
-  private void executeCreateTimeseriesByDeviceStatement(
-      CreateTimeSeriesByDeviceStatement statement) {
-    long queryId = SessionManager.getInstance().requestQueryId(false);
-    ExecutionResult executionResult =
-        coordinator.execute(statement, queryId, null, "", partitionFetcher, 
this);
-    // TODO: throw exception
     try {
-      int statusCode = executionResult.status.getCode();
-      if (statusCode == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        return;
-      }
-
-      for (TSStatus subStatus : executionResult.status.subStatus) {
-        if (subStatus.code != 
TSStatusCode.PATH_ALREADY_EXIST_ERROR.getStatusCode()) {
-          throw new RuntimeException(
-              "cannot auto create schema, status is: " + 
executionResult.status);
+      if (isAligned) {
+        CreateAlignedTimeSeriesPlan createAlignedTimeSeriesPlan = new 
CreateAlignedTimeSeriesPlan();
+        createAlignedTimeSeriesPlan.setPrefixPath(devicePath);
+        createAlignedTimeSeriesPlan.setMeasurements(measurements);
+        createAlignedTimeSeriesPlan.setDataTypes(tsDataTypes);
+        List<TSEncoding> encodings = new ArrayList<>();
+        List<CompressionType> compressors = new ArrayList<>();
+        for (TSDataType dataType : tsDataTypes) {
+          encodings.add(getDefaultEncoding(dataType));
+          
compressors.add(TSFileDescriptor.getInstance().getConfig().getCompressor());
+        }
+        createAlignedTimeSeriesPlan.setEncodings(encodings);
+        createAlignedTimeSeriesPlan.setCompressors(compressors);
+        SchemaRegionId schemaRegionId =
+            
localConfigNode.getBelongedSchemaRegionIdWithAutoCreate(devicePath);
+        ISchemaRegion schemaRegion = 
schemaEngine.getSchemaRegion(schemaRegionId);
+        schemaRegion.createAlignedTimeSeries(createAlignedTimeSeriesPlan);
+      } else {
+        for (int i = 0; i < measurements.size(); i++) {
+          CreateTimeSeriesPlan createTimeSeriesPlan = new 
CreateTimeSeriesPlan();
+          createTimeSeriesPlan.setPath(
+              new PartialPath(devicePath.getFullPath(), measurements.get(i)));
+          createTimeSeriesPlan.setDataType(tsDataTypes.get(i));
+          
createTimeSeriesPlan.setEncoding(getDefaultEncoding(tsDataTypes.get(i)));
+          createTimeSeriesPlan.setCompressor(
+              TSFileDescriptor.getInstance().getConfig().getCompressor());
+          SchemaRegionId schemaRegionId =
+              
localConfigNode.getBelongedSchemaRegionIdWithAutoCreate(devicePath);
+          ISchemaRegion schemaRegion = 
schemaEngine.getSchemaRegion(schemaRegionId);
+          schemaRegion.createTimeseries(createTimeSeriesPlan, -1);
         }
       }
-    } finally {
-      coordinator.getQueryExecution(queryId).stopAndCleanup();
+    } catch (MetadataException e) {
+      throw new RuntimeException("cannot auto create schema ", e);
     }
   }
 }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/plan/StandaloneCoordinatorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/plan/StandaloneCoordinatorTest.java
index 46ceab5ef3..9aeed28fcc 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/plan/StandaloneCoordinatorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/plan/StandaloneCoordinatorTest.java
@@ -77,6 +77,7 @@ public class StandaloneCoordinatorTest {
   @After
   public void tearDown() throws Exception {
     configNode.clear();
+    WALManager.getInstance().clear();
     WALManager.getInstance().stop();
     StorageEngineV2.getInstance().stop();
     FlushManager.getInstance().stop();
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/plan/scheduler/StandaloneSchedulerTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/plan/scheduler/StandaloneSchedulerTest.java
index ee16c73244..ce58dd7e41 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/plan/scheduler/StandaloneSchedulerTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/plan/scheduler/StandaloneSchedulerTest.java
@@ -155,9 +155,15 @@ public class StandaloneSchedulerTest {
             executor,
             null,
             null);
-    standaloneScheduler.start();
-
-    Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    try {
+      standaloneScheduler.start();
+      Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    } catch (Exception e) {
+      e.printStackTrace();
+      Assert.fail();
+    } finally {
+      standaloneScheduler.stop();
+    }
   }
 
   @Test
@@ -248,9 +254,15 @@ public class StandaloneSchedulerTest {
             executor,
             null,
             null);
-    standaloneScheduler.start();
-
-    Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    try {
+      standaloneScheduler.start();
+      Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    } catch (Exception e) {
+      e.printStackTrace();
+      Assert.fail();
+    } finally {
+      standaloneScheduler.stop();
+    }
   }
 
   @Test
@@ -351,9 +363,15 @@ public class StandaloneSchedulerTest {
             executor,
             null,
             null);
-    standaloneScheduler.start();
-
-    Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    try {
+      standaloneScheduler.start();
+      Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    } catch (Exception e) {
+      e.printStackTrace();
+      Assert.fail();
+    } finally {
+      standaloneScheduler.stop();
+    }
   }
 
   @Test
@@ -393,8 +411,15 @@ public class StandaloneSchedulerTest {
             executor,
             null,
             null);
-    standaloneScheduler.start();
-    Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    try {
+      standaloneScheduler.start();
+      Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    } catch (Exception e) {
+      e.printStackTrace();
+      Assert.fail();
+    } finally {
+      standaloneScheduler.stop();
+    }
   }
 
   @Test
@@ -463,8 +488,15 @@ public class StandaloneSchedulerTest {
             executor,
             null,
             null);
-    standaloneScheduler.start();
-    Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    try {
+      standaloneScheduler.start();
+      Assert.assertEquals(QueryState.FINISHED, stateMachine.getState());
+    } catch (Exception e) {
+      e.printStackTrace();
+      Assert.fail();
+    } finally {
+      standaloneScheduler.stop();
+    }
   }
 
   private TRegionReplicaSet genRegionReplicaSet(TConsensusGroupType type) {

Reply via email to