This is an automated email from the ASF dual-hosted git repository. caogaofei pushed a commit to branch beyyes/multi_devices_fe in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 54826e2734456d143f997c49893b2d828376c56b Author: Beyyes <[email protected]> AuthorDate: Wed Oct 25 14:39:13 2023 +0800 add log for coordinator, fix code smell --- .../apache/iotdb/db/queryengine/plan/Coordinator.java | 4 +++- .../db/queryengine/plan/analyze/AnalyzeVisitor.java | 4 ++++ .../db/queryengine/plan/execution/QueryExecution.java | 3 +++ .../plan/distribution/DistributionPlannerCycleTest.java | 9 +++++---- .../db/queryengine/plan/plan/distribution/Util2.java | 16 +++++++++------- 5 files changed, 24 insertions(+), 12 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/Coordinator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/Coordinator.java index cf4eb213752..a3c58315f07 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/Coordinator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/Coordinator.java @@ -158,7 +158,9 @@ public class Coordinator { queryContext.setTimeOut(Long.MAX_VALUE); } execution.start(); - + LOGGER.warn( + "========= Consume time in Coordinator.execute: {}ms", + System.currentTimeMillis() - startTime); return execution.getStatus(); } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/AnalyzeVisitor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/AnalyzeVisitor.java index 6b47b26ad0e..fc8f341d8d0 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/AnalyzeVisitor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/AnalyzeVisitor.java @@ -274,6 +274,10 @@ public class AnalyzeVisitor extends StatementVisitor<Analysis, MPPQueryContext> List<Pair<Expression, String>> outputExpressions; if (queryStatement.isAlignByDevice()) { + // TODO + // 并行 最上层比如加sort + // fe 线程管理 + // List<PartialPath> deviceList = analyzeFrom(queryStatement, schemaTree); if (canPushDownLimitOffsetInGroupByTimeForDevice(queryStatement)) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/QueryExecution.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/QueryExecution.java index a002f99cb29..1fd4652a0bb 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/QueryExecution.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/QueryExecution.java @@ -239,6 +239,9 @@ public class QueryExecution implements IQueryExecution { if (context.getQueryType() == QueryType.WRITE && analysis.isFailed()) { stateMachine.transitionToFailed(analysis.getFailStatus()); } + logger.warn( + "~~~~ Consume time in doLogicalPlan+doDistributionPlan: {}ns", + System.nanoTime() - startTime); } private void checkTimeOutForQuery() { diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/plan/distribution/DistributionPlannerCycleTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/plan/distribution/DistributionPlannerCycleTest.java index 8961c1e207f..8af58bb3c5a 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/plan/distribution/DistributionPlannerCycleTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/plan/distribution/DistributionPlannerCycleTest.java @@ -40,7 +40,7 @@ public class DistributionPlannerCycleTest { // Query sql: `select * from root.sg.d1,root.sg.d2` // root.sg.d1 has 2 SeriesScanNodes, root.sg.d2 has 3 SeriesScanNodes. // - // ------------------------------------------------------------------------------------------------ + // ----------------------------------------------------------------------------------------- // Note: d1.s1[1] means a SeriesScanNode with target series d1.s1 and its data region is 1 // // IdentityNode @@ -51,7 +51,7 @@ public class DistributionPlannerCycleTest { // TimeJoinNode // / \ \ // d2.s1[2] d2.s2[2] d2.s3[2] - // ------------------------------------------------------------------------------------------------ + // ------------------------------------------------------------------------------------------ @Test public void timeJoinNodeTest() { QueryId queryId = new QueryId("test"); @@ -67,12 +67,13 @@ public class DistributionPlannerCycleTest { assertEquals(2, plan.getInstances().size()); PlanNode firstNode = plan.getInstances().get(0).getFragment().getPlanNodeTree().getChildren().get(0); - PlanNode secondNode = - plan.getInstances().get(1).getFragment().getPlanNodeTree().getChildren().get(0); assertEquals(3, firstNode.getChildren().size()); assertTrue(firstNode.getChildren().get(0) instanceof SeriesScanNode); assertTrue(firstNode.getChildren().get(1) instanceof SeriesScanNode); assertTrue(firstNode.getChildren().get(2) instanceof ExchangeNode); + + PlanNode secondNode = + plan.getInstances().get(1).getFragment().getPlanNodeTree().getChildren().get(0); assertEquals(3, secondNode.getChildren().size()); assertTrue(secondNode.getChildren().get(0) instanceof SeriesScanNode); assertTrue(secondNode.getChildren().get(1) instanceof SeriesScanNode); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/plan/distribution/Util2.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/plan/distribution/Util2.java index 5124c61c28b..89a072b0792 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/plan/distribution/Util2.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/plan/distribution/Util2.java @@ -71,6 +71,10 @@ import java.util.Map; public class Util2 { public static final Analysis ANALYSIS = constructAnalysis(); + private static final String device1 = "root.sg.d1"; + private static final String device2 = "root.sg.d2"; + private static final String device3 = "root.sg.d3"; + public static Analysis constructAnalysis() { TRegionReplicaSet dataRegion1 = new TRegionReplicaSet( @@ -90,23 +94,21 @@ public class Util2 { Map<TTimePartitionSlot, List<TRegionReplicaSet>> d2DataRegionMap = new HashMap<>(); d2DataRegionMap.put(new TTimePartitionSlot(), d2DataRegions); - DataPartition dataPartition = - new DataPartition( - IoTDBDescriptor.getInstance().getConfig().getSeriesPartitionExecutorClass(), - IoTDBDescriptor.getInstance().getConfig().getSeriesPartitionSlotNum()); Map<String, Map<TSeriesPartitionSlot, Map<TTimePartitionSlot, List<TRegionReplicaSet>>>> dataPartitionMap = new HashMap<>(); Map<TSeriesPartitionSlot, Map<TTimePartitionSlot, List<TRegionReplicaSet>>> sgPartitionMap = new HashMap<>(); - String device1 = "root.sg.d1"; - String device2 = "root.sg.d2"; - String device3 = "root.sg.d3"; + SeriesPartitionExecutor executor = SeriesPartitionExecutor.getSeriesPartitionExecutor( IoTDBDescriptor.getInstance().getConfig().getSeriesPartitionExecutorClass(), IoTDBDescriptor.getInstance().getConfig().getSeriesPartitionSlotNum()); sgPartitionMap.put(executor.getSeriesPartitionSlot(device1), d1DataRegionMap); sgPartitionMap.put(executor.getSeriesPartitionSlot(device2), d2DataRegionMap); + DataPartition dataPartition = + new DataPartition( + IoTDBDescriptor.getInstance().getConfig().getSeriesPartitionExecutorClass(), + IoTDBDescriptor.getInstance().getConfig().getSeriesPartitionSlotNum()); dataPartitionMap.put("root.sg", sgPartitionMap); dataPartition.setDataPartitionMap(dataPartitionMap);
