This is an automated email from the ASF dual-hosted git repository. xiangweiwei pushed a commit to branch addAggregationIT in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 215fe41dd9133a0516ffdfdff6df4a20367b7f76 Author: Alima777 <[email protected]> AuthorDate: Fri Jun 17 16:21:58 2022 +0800 add group by level IT --- .../itbase/runtime/ParallelRequestDelegate.java | 5 +- .../itbase/runtime/SerialRequestDelegate.java | 3 +- .../it}/aggregation/IoTDBAggregationByLevelIT.java | 322 ++++++++++----------- .../apache/iotdb/db/mpp/plan/analyze/Analyzer.java | 19 +- .../db/mpp/plan/analyze/ExpressionAnalyzer.java | 1 + .../mpp/plan/analyze/GroupByLevelController.java | 13 +- .../statement/component/GroupByLevelComponent.java | 13 + 7 files changed, 201 insertions(+), 175 deletions(-) diff --git a/integration-test/src/main/java/org/apache/iotdb/itbase/runtime/ParallelRequestDelegate.java b/integration-test/src/main/java/org/apache/iotdb/itbase/runtime/ParallelRequestDelegate.java index 32d9803032..4ef07dc14a 100644 --- a/integration-test/src/main/java/org/apache/iotdb/itbase/runtime/ParallelRequestDelegate.java +++ b/integration-test/src/main/java/org/apache/iotdb/itbase/runtime/ParallelRequestDelegate.java @@ -52,7 +52,10 @@ public class ParallelRequestDelegate<T> extends RequestDelegate<T> { resultFutures.get(j).cancel(true); } throw new SQLException( - String.format("Waiting for query results of %s failed", getEndpoints().get(i)), e); + String.format( + "Waiting for query results of %s failed: %s", + getEndpoints().get(i), e.getMessage()), + e); } } return results; diff --git a/integration-test/src/main/java/org/apache/iotdb/itbase/runtime/SerialRequestDelegate.java b/integration-test/src/main/java/org/apache/iotdb/itbase/runtime/SerialRequestDelegate.java index 10e53b62e6..ade0f1ed67 100644 --- a/integration-test/src/main/java/org/apache/iotdb/itbase/runtime/SerialRequestDelegate.java +++ b/integration-test/src/main/java/org/apache/iotdb/itbase/runtime/SerialRequestDelegate.java @@ -39,7 +39,8 @@ public class SerialRequestDelegate<T> extends RequestDelegate<T> { try { results.add(getRequests().get(i).call()); } catch (Exception e) { - throw new SQLException(String.format("Request %s error.", getEndpoints().get(i)), e); + throw new SQLException( + String.format("Request %s error: %s", getEndpoints().get(i), e.getMessage()), e); } } return results; diff --git a/integration/src/test/java/org/apache/iotdb/db/integration/aggregation/IoTDBAggregationByLevelIT.java b/integration-test/src/test/java/org/apache/iotdb/db/it/aggregation/IoTDBAggregationByLevelIT.java similarity index 61% rename from integration/src/test/java/org/apache/iotdb/db/integration/aggregation/IoTDBAggregationByLevelIT.java rename to integration-test/src/test/java/org/apache/iotdb/db/it/aggregation/IoTDBAggregationByLevelIT.java index badd79a9af..bb8e0aefb9 100644 --- a/integration/src/test/java/org/apache/iotdb/db/integration/aggregation/IoTDBAggregationByLevelIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/db/it/aggregation/IoTDBAggregationByLevelIT.java @@ -16,17 +16,16 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.iotdb.db.integration.aggregation; +package org.apache.iotdb.db.it.aggregation; -import org.apache.iotdb.db.constant.TestConstant; -import org.apache.iotdb.db.qp.Planner; -import org.apache.iotdb.db.qp.logical.crud.AggregationQueryOperator; -import org.apache.iotdb.integration.env.EnvFactory; -import org.apache.iotdb.itbase.category.LocalStandaloneTest; +import org.apache.iotdb.it.env.EnvFactory; +import org.apache.iotdb.itbase.category.ClusterIT; +import org.apache.iotdb.itbase.category.LocalStandaloneIT; import org.junit.AfterClass; import org.junit.Assert; import org.junit.BeforeClass; +import org.junit.Ignore; import org.junit.Test; import org.junit.experimental.categories.Category; @@ -34,12 +33,19 @@ import java.sql.Connection; import java.sql.ResultSet; import java.sql.Statement; +import static org.apache.iotdb.itbase.constant.TestConstant.TIMESTAMP_STR; +import static org.apache.iotdb.itbase.constant.TestConstant.avg; +import static org.apache.iotdb.itbase.constant.TestConstant.count; +import static org.apache.iotdb.itbase.constant.TestConstant.lastValue; +import static org.apache.iotdb.itbase.constant.TestConstant.maxTime; +import static org.apache.iotdb.itbase.constant.TestConstant.maxValue; +import static org.apache.iotdb.itbase.constant.TestConstant.minTime; +import static org.apache.iotdb.itbase.constant.TestConstant.sum; import static org.junit.Assert.fail; -@Category({LocalStandaloneTest.class}) +@Category({LocalStandaloneIT.class, ClusterIT.class}) public class IoTDBAggregationByLevelIT { - private Planner planner = new Planner(); private static final String[] dataSet = new String[] { "SET STORAGE GROUP TO root.sg1", @@ -83,41 +89,41 @@ public class IoTDBAggregationByLevelIT { Statement statement = connection.createStatement()) { // Here we duplicate the column to test the bug // https://issues.apache.org/jira/browse/IOTDB-2088 - statement.execute( - "select sum(temperature), sum(temperature) from root.sg1.* GROUP BY level=1"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select sum(temperature), sum(temperature) from root.sg1.* GROUP BY level=1")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.sum("root.sg1.*.temperature")); + String ans = resultSet.getString(sum("root.sg1.*.temperature")); Assert.assertEquals(retArray[cnt], Double.parseDouble(ans), DOUBLE_PRECISION); cnt++; } } - statement.execute("select sum(temperature) from root.sg2.* GROUP BY level=1"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select sum(temperature) from root.sg2.* GROUP BY level=1")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.sum("root.sg2.*.temperature")); + String ans = resultSet.getString(sum("root.sg2.*.temperature")); Assert.assertEquals(retArray[cnt], Double.parseDouble(ans), DOUBLE_PRECISION); cnt++; } } - statement.execute("select sum(temperature) from root.*.* GROUP BY level=0"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select sum(temperature) from root.*.* GROUP BY level=0")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.sum("root.*.*.temperature")); + String ans = resultSet.getString(sum("root.*.*.temperature")); Assert.assertEquals(retArray[cnt], Double.parseDouble(ans), DOUBLE_PRECISION); cnt++; } } - statement.execute("select sum(temperature) from root.sg1.* GROUP BY level=1,2"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select sum(temperature) from root.sg1.* GROUP BY level=1,2")) { while (resultSet.next()) { - String ans1 = resultSet.getString(TestConstant.sum("root.sg1.d1.temperature")); - String ans2 = resultSet.getString(TestConstant.sum("root.sg1.d2.temperature")); + String ans1 = resultSet.getString(sum("root.sg1.d1.temperature")); + String ans2 = resultSet.getString(sum("root.sg1.d2.temperature")); Assert.assertEquals(retArray[cnt++], Double.parseDouble(ans1), DOUBLE_PRECISION); Assert.assertEquals(retArray[cnt++], Double.parseDouble(ans2), DOUBLE_PRECISION); } @@ -131,40 +137,40 @@ public class IoTDBAggregationByLevelIT { double[] retArray = new double[] {48.682d, 95.115d, 69.319d, 45.915d, 50.527d}; try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute("select avg(temperature) from root.sg1.* GROUP BY level=1"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select avg(temperature) from root.sg1.* GROUP BY level=1")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.avg("root.sg1.*.temperature")); + String ans = resultSet.getString(avg("root.sg1.*.temperature")); Assert.assertEquals(retArray[cnt], Double.parseDouble(ans), DOUBLE_PRECISION); cnt++; } } - statement.execute("select avg(temperature) from root.sg2.* GROUP BY level=1"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select avg(temperature) from root.sg2.* GROUP BY level=1")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.avg("root.sg2.*.temperature")); + String ans = resultSet.getString(avg("root.sg2.*.temperature")); Assert.assertEquals(retArray[cnt], Double.parseDouble(ans), DOUBLE_PRECISION); cnt++; } } - statement.execute("select avg(temperature) from root.*.* GROUP BY level=0"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select avg(temperature) from root.*.* GROUP BY level=0")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.avg("root.*.*.temperature")); + String ans = resultSet.getString(avg("root.*.*.temperature")); Assert.assertEquals(retArray[cnt], Double.parseDouble(ans), DOUBLE_PRECISION); cnt++; } } - statement.execute("select avg(temperature) from root.sg1.* GROUP BY level=1, 2"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select avg(temperature) from root.sg1.* GROUP BY level=1, 2")) { while (resultSet.next()) { - String ans1 = resultSet.getString(TestConstant.avg("root.sg1.d1.temperature")); - String ans2 = resultSet.getString(TestConstant.avg("root.sg1.d2.temperature")); + String ans1 = resultSet.getString(avg("root.sg1.d1.temperature")); + String ans2 = resultSet.getString(avg("root.sg1.d2.temperature")); Assert.assertEquals(retArray[cnt++], Double.parseDouble(ans1), DOUBLE_PRECISION); Assert.assertEquals(retArray[cnt++], Double.parseDouble(ans2), DOUBLE_PRECISION); } @@ -178,51 +184,51 @@ public class IoTDBAggregationByLevelIT { String[] retArray = new String[] {"5,3,100,200", "600,700,2,3", "600,700,500"}; try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( - "select count(status), min_time(temperature) from root.*.* GROUP BY level=2"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(status), min_time(temperature) from root.*.* GROUP BY level=2")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.count("root.*.d1.status")) + resultSet.getString(count("root.*.d1.status")) + "," - + resultSet.getString(TestConstant.count("root.*.d2.status")) + + resultSet.getString(count("root.*.d2.status")) + "," - + resultSet.getString(TestConstant.minTime("root.*.d1.temperature")) + + resultSet.getString(minTime("root.*.d1.temperature")) + "," - + resultSet.getString(TestConstant.minTime("root.*.d2.temperature")); + + resultSet.getString(minTime("root.*.d2.temperature")); Assert.assertEquals(retArray[cnt], ans); cnt++; } } - statement.execute( - "select max_time(status), count(temperature) from root.sg1.* GROUP BY level=2"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select max_time(status), count(temperature) from root.sg1.* GROUP BY level=2")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.maxTime("root.*.d1.status")) + resultSet.getString(maxTime("root.*.d1.status")) + "," - + resultSet.getString(TestConstant.maxTime("root.*.d2.status")) + + resultSet.getString(maxTime("root.*.d2.status")) + "," - + resultSet.getString(TestConstant.count("root.*.d1.temperature")) + + resultSet.getString(count("root.*.d1.temperature")) + "," - + resultSet.getString(TestConstant.count("root.*.d2.temperature")); + + resultSet.getString(count("root.*.d2.temperature")); Assert.assertEquals(retArray[cnt], ans); cnt++; } } - statement.execute("select max_time(status) from root.*.* GROUP BY level=1, 2"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select max_time(status) from root.*.* GROUP BY level=1, 2")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.maxTime("root.sg1.d1.status")) + resultSet.getString(maxTime("root.sg1.d1.status")) + "," - + resultSet.getString(TestConstant.maxTime("root.sg1.d2.status")) + + resultSet.getString(maxTime("root.sg1.d2.status")) + "," - + resultSet.getString(TestConstant.maxTime("root.sg2.d1.status")); + + resultSet.getString(maxTime("root.sg2.d1.status")); Assert.assertEquals(retArray[cnt], ans); cnt++; } @@ -239,33 +245,33 @@ public class IoTDBAggregationByLevelIT { }; try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( - "select last_value(temperature), max_value(temperature) from root.*.* GROUP BY level=0"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select last_value(temperature), max_value(temperature) from root.*.* GROUP BY level=0")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.lastValue("root.*.*.temperature")) + resultSet.getString(lastValue("root.*.*.temperature")) + "," - + resultSet.getString(TestConstant.maxValue("root.*.*.temperature")); + + resultSet.getString(maxValue("root.*.*.temperature")); Assert.assertEquals(retArray[cnt], ans); cnt++; } } - statement.execute( - "select last_value(temperature), max_value(temperature) from root.sg1.* GROUP BY level=2"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select last_value(temperature), max_value(temperature) from root.sg1.* GROUP BY level=2")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.lastValue("root.*.d1.temperature")) + resultSet.getString(lastValue("root.*.d1.temperature")) + "," - + resultSet.getString(TestConstant.lastValue("root.*.d2.temperature")) + + resultSet.getString(lastValue("root.*.d2.temperature")) + "," - + resultSet.getString(TestConstant.maxValue("root.*.d1.temperature")) + + resultSet.getString(maxValue("root.*.d1.temperature")) + "," - + resultSet.getString(TestConstant.maxValue("root.*.d2.temperature")); + + resultSet.getString(maxValue("root.*.d2.temperature")); Assert.assertEquals(retArray[cnt], ans); cnt++; } @@ -279,32 +285,30 @@ public class IoTDBAggregationByLevelIT { String[] retArray = new String[] {"17", "17", "8"}; try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute("select count(*) from root.*.* GROUP BY level=0"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select count(*) from root.*.* GROUP BY level=0")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.count("root.*.*.*")); + String ans = resultSet.getString(count("root.*.*.*")); Assert.assertEquals(retArray[cnt], ans); cnt++; } } - statement.execute("select count(**) from root GROUP BY level=0"); - - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select count(**) from root GROUP BY level=0")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.count("root.*.*.*")); + String ans = resultSet.getString(count("root.*.*.*")); Assert.assertEquals(retArray[cnt], ans); cnt++; } } - statement.execute("select count(status) from root.*.* GROUP BY level=0"); - - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select count(status) from root.*.* GROUP BY level=0")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.count("root.*.*.status")); + String ans = resultSet.getString(count("root.*.*.status")); Assert.assertEquals(retArray[cnt], ans); cnt++; } @@ -322,11 +326,11 @@ public class IoTDBAggregationByLevelIT { String[] retArray = new String[] {"5", "5", "5"}; try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( - "select count(temperature) as ct from root.sg1.d1, root.sg1.d2 GROUP BY level=1"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature) as ct from root.sg1.d1, root.sg1.d2 GROUP BY level=1")) { while (resultSet.next()) { String ans = resultSet.getString("ct"); Assert.assertEquals(retArray[cnt], ans); @@ -334,9 +338,10 @@ public class IoTDBAggregationByLevelIT { } } - statement.execute("select count(temperature) as ct from root.sg1.* GROUP BY level=1"); cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature) as ct from root.sg1.* GROUP BY level=1")) { while (resultSet.next()) { String ans = resultSet.getString("ct"); Assert.assertEquals(retArray[cnt], ans); @@ -345,8 +350,8 @@ public class IoTDBAggregationByLevelIT { } // root.sg1.d1.* -> [root.sg1.d1.status, root.sg1.d1.temperature] -> root.*.*.* -> ct - statement.execute("select count(*) as ct from root.sg1.d1 GROUP BY level=0"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery("select count(*) as ct from root.sg1.d1 GROUP BY level=0")) { while (resultSet.next()) { String ans = resultSet.getString("ct"); Assert.assertEquals(retArray[cnt], ans); @@ -365,7 +370,7 @@ public class IoTDBAggregationByLevelIT { public void groupByLevelWithAliasFailTest() throws Exception { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute("select count(temperature) as ct from root.sg1.* GROUP BY level=2"); + statement.executeQuery("select count(temperature) as ct from root.sg1.* GROUP BY level=2"); fail("No exception thrown"); } catch (Exception e) { Assert.assertTrue(e.getMessage().contains("can only be matched with one")); @@ -377,10 +382,11 @@ public class IoTDBAggregationByLevelIT { public void groupByLevelWithAliasFailTest2() throws Exception { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( + statement.executeQuery( "select count(temperature) as ct from root.sg1.d1, root.sg2.d2 GROUP BY level=2"); fail("No exception thrown"); } catch (Exception e) { + System.out.println(e.getMessage()); Assert.assertTrue(e.getMessage().contains("can only be matched with one")); } } @@ -390,7 +396,7 @@ public class IoTDBAggregationByLevelIT { public void groupByLevelWithAliasFailTest3() throws Exception { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( + statement.executeQuery( "select count(temperature) as ct, count(temperature) as ct2 from root.sg1.d1 GROUP BY level=2"); fail("No exception thrown"); } catch (Exception e) { @@ -404,11 +410,11 @@ public class IoTDBAggregationByLevelIT { String[] retArray2 = new String[] {"0,0", "100,1", "200,2", "300,0", "400,0", "500,0"}; try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( - "select count(temperature) as ct from root.sg1.d1, root.sg1.d2 GROUP BY ([0, 600), 100ms), level=1"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature) as ct from root.sg1.d1, root.sg1.d2 GROUP BY ([0, 600), 100ms), level=1")) { while (resultSet.next()) { String ans = resultSet.getString("Time") + "," + resultSet.getString("ct"); Assert.assertEquals(retArray[cnt], ans); @@ -416,11 +422,10 @@ public class IoTDBAggregationByLevelIT { } } - statement.execute( - "select count(temperature) as ct from root.sg1.* GROUP BY ([0, 600), 100ms), level=1"); - cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature) as ct from root.sg1.* GROUP BY ([0, 600), 100ms), level=1")) { while (resultSet.next()) { String ans = resultSet.getString("Time") + "," + resultSet.getString("ct"); Assert.assertEquals(retArray[cnt], ans); @@ -430,9 +435,9 @@ public class IoTDBAggregationByLevelIT { cnt = 0; // root.sg1.d1.* -> [root.sg1.d1.status, root.sg1.d1.temperature] -> root.*.*.* -> ct - statement.execute( - "select count(*) as ct from root.sg1.d1 GROUP BY ([0, 600), 100ms), level=1"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(*) as ct from root.sg1.d1 GROUP BY ([0, 600), 100ms), level=1")) { while (resultSet.next()) { String ans = resultSet.getString("Time") + "," + resultSet.getString("ct"); Assert.assertEquals(retArray2[cnt], ans); @@ -451,7 +456,7 @@ public class IoTDBAggregationByLevelIT { public void groupByLevelWithAliasWithTimeIntervalFailTest() throws Exception { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( + statement.executeQuery( "select count(temperature) as ct from root.sg1.* GROUP BY ([0, 600), 100ms), level=2"); fail(); } catch (Exception e) { @@ -464,7 +469,7 @@ public class IoTDBAggregationByLevelIT { public void groupByLevelWithAliasWithTimeIntervalFailTest2() throws Exception { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( + statement.executeQuery( "select count(temperature) as ct from root.sg1.d1, root.sg1.d2 GROUP BY ([0, 600), 100ms), level=2"); fail(); } catch (Exception e) { @@ -477,41 +482,39 @@ public class IoTDBAggregationByLevelIT { String[] retArray = new String[] {"5,4", "4,6", "3"}; try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( - "select count(temperature), count(status) from root.*.* GROUP BY level=1 slimit 2"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature), count(status) from root.*.* GROUP BY level=1 slimit 2")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.count("root.sg1.*.temperature")) + resultSet.getString(count("root.sg1.*.temperature")) + "," - + resultSet.getString(TestConstant.count("root.sg2.*.temperature")); + + resultSet.getString(count("root.sg2.*.temperature")); Assert.assertEquals(retArray[cnt], ans); cnt++; } } - statement.execute( - "select count(temperature), count(status) from root.*.* GROUP BY level=1 slimit 2 soffset 1"); - - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature), count(status) from root.*.* GROUP BY level=1 slimit 2 soffset 1")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.count("root.sg2.*.temperature")) + resultSet.getString(count("root.sg2.*.temperature")) + "," - + resultSet.getString(TestConstant.count("root.sg1.*.status")); + + resultSet.getString(count("root.sg1.*.status")); Assert.assertEquals(retArray[cnt], ans); cnt++; } } - statement.execute( - "select count(temperature), count(status) from root.*.* GROUP BY level=1,2 slimit 1 soffset 4"); - - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature), count(status) from root.*.* GROUP BY level=1,2 slimit 1 soffset 4")) { while (resultSet.next()) { - String ans = resultSet.getString(TestConstant.count("root.sg1.d1.status")); + String ans = resultSet.getString(count("root.sg1.d1.status")); Assert.assertEquals(retArray[cnt], ans); cnt++; } @@ -537,30 +540,30 @@ public class IoTDBAggregationByLevelIT { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( - "select sum(temperature) from root.sg2.* GROUP BY ([0, 600), 100ms), level=1"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select sum(temperature) from root.sg2.* GROUP BY ([0, 600), 100ms), level=1")) { while (resultSet.next()) { - String ans = "" + resultSet.getString(TestConstant.sum("root.sg2.*.temperature")); + String ans = "" + resultSet.getString(sum("root.sg2.*.temperature")); Assert.assertEquals(retArray1[cnt], ans); cnt++; } } cnt = 0; - statement.execute( - "select max_time(temperature), avg(temperature) from root.*.* GROUP BY ([0, 600), 100ms), level=1"); - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select max_time(temperature), avg(temperature) from root.*.* GROUP BY ([0, 600), 100ms), level=1")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.maxTime("root.sg1.*.temperature")) + resultSet.getString(maxTime("root.sg1.*.temperature")) + "," - + resultSet.getString(TestConstant.maxTime("root.sg2.*.temperature")) + + resultSet.getString(maxTime("root.sg2.*.temperature")) + "," - + resultSet.getString(TestConstant.avg("root.sg1.*.temperature")) + + resultSet.getString(avg("root.sg1.*.temperature")) + "," - + resultSet.getString(TestConstant.avg("root.sg2.*.temperature")); + + resultSet.getString(avg("root.sg2.*.temperature")); Assert.assertEquals(retArray2[cnt], ans); cnt++; } @@ -577,16 +580,17 @@ public class IoTDBAggregationByLevelIT { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( - "select sum(temperature) from root.sg2.* GROUP BY ([0, 600), 100ms), level=0,1"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select sum(temperature) from root.sg2.* GROUP BY ([0, 600), 100ms), level=0,1")) { while (resultSet.next()) { - String ans = "" + resultSet.getString(TestConstant.sum("root.sg2.*.temperature")); + String ans = "" + resultSet.getString(sum("root.sg2.*.temperature")); Assert.assertEquals(retArray1[cnt], ans); cnt++; } } + Assert.assertEquals(6, cnt); } } @@ -600,33 +604,33 @@ public class IoTDBAggregationByLevelIT { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute( - "select count(temperature) from root.*.* GROUP BY ([0, 600), 100ms), level=1 slimit 2"); int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature) from root.*.* GROUP BY ([0, 600), 100ms), level=1 slimit 2")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.TIMESTAMP_STR) + resultSet.getString(TIMESTAMP_STR) + "," - + resultSet.getString(TestConstant.count("root.sg1.*.temperature")) + + resultSet.getString(count("root.sg1.*.temperature")) + "," - + resultSet.getString(TestConstant.count("root.sg2.*.temperature")); + + resultSet.getString(count("root.sg2.*.temperature")); Assert.assertEquals(retArray[cnt], ans); cnt++; } } - statement.execute( - "select count(temperature), count(status) from root.*.* GROUP BY ([0, 600), 100ms), level=1 slimit 2 soffset 1"); cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { + try (ResultSet resultSet = + statement.executeQuery( + "select count(temperature), count(status) from root.*.* GROUP BY ([0, 600), 100ms), level=1 slimit 2 soffset 1")) { while (resultSet.next()) { String ans = - resultSet.getString(TestConstant.TIMESTAMP_STR) + resultSet.getString(TIMESTAMP_STR) + "," - + resultSet.getString(TestConstant.count("root.sg2.*.temperature")) + + resultSet.getString(count("root.sg2.*.temperature")) + "," - + resultSet.getString(TestConstant.count("root.sg1.*.status")); + + resultSet.getString(count("root.sg1.*.status")); Assert.assertEquals(retArray2[cnt], ans); cnt++; } @@ -636,28 +640,12 @@ public class IoTDBAggregationByLevelIT { @Test public void mismatchedFuncGroupByLevelTest() throws Exception { - String[] retArray = - new String[] { - "true", "3", - }; try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute("select last_value(status) from root.*.* GROUP BY level=0"); - - int cnt = 0; - try (ResultSet resultSet = statement.getResultSet()) { - while (resultSet.next()) { - String ans = resultSet.getString(1); - Assert.assertEquals(retArray[cnt], ans); - cnt++; - } - } - - try { - planner.parseSQLToPhysicalPlan("select avg(status) from root.sg2.* GROUP BY level=1"); - } catch (Exception e) { - Assert.assertEquals("Aggregate among unmatched data types", e.getMessage()); - } + statement.executeQuery("select last_value(status) from root.*.* GROUP BY level=0"); + fail(); + } catch (Exception e) { + Assert.assertTrue(e.getMessage().contains("the data types of the same output column")); } } @@ -665,16 +653,20 @@ public class IoTDBAggregationByLevelIT { * Test group by level without aggregation function used in select clause. The expected situation * is throwing an exception. */ + // TODO + @Ignore @Test public void TestGroupByLevelWithoutAggregationFunc() { try (Connection connection = EnvFactory.getEnv().getConnection(); Statement statement = connection.createStatement()) { - statement.execute("select temperature from root.sg1.* group by level = 2"); - + statement.executeQuery("select temperature from root.sg1.* group by level = 2"); fail("No expected exception thrown"); } catch (Exception e) { - Assert.assertTrue(e.getMessage().contains(AggregationQueryOperator.ERROR_MESSAGE1)); + Assert.assertTrue( + e.getMessage() + .contains( + "Common queries and aggregated queries are not allowed to appear at the same time")); } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java index d1e3c4cdf8..93cc167e81 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java @@ -40,6 +40,7 @@ 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.expression.Expression; import org.apache.iotdb.db.mpp.plan.expression.leaf.TimeSeriesOperand; +import org.apache.iotdb.db.mpp.plan.expression.multi.FunctionExpression; import org.apache.iotdb.db.mpp.plan.planner.plan.parameter.FillDescriptor; import org.apache.iotdb.db.mpp.plan.planner.plan.parameter.FilterNullParameter; import org.apache.iotdb.db.mpp.plan.planner.plan.parameter.GroupByTimeParameter; @@ -430,7 +431,8 @@ public class Analyzer { analysis.setDataPartitionInfo(dataPartition); } catch (StatementAnalyzeException e) { logger.error("Meet error when analyzing the query statement: ", e); - throw new StatementAnalyzeException("Meet error when analyzing the query statement"); + throw new StatementAnalyzeException( + "Meet error when analyzing the query statement: " + e.getMessage()); } return analysis; } @@ -448,7 +450,7 @@ public class Analyzer { boolean hasAlias = resultColumn.hasAlias(); List<Expression> resultExpressions = ExpressionAnalyzer.removeWildcardInExpression(resultColumn.getExpression(), schemaTree); - if (hasAlias && resultExpressions.size() > 1) { + if (hasAlias && !queryStatement.isGroupByLevel() && resultExpressions.size() > 1) { throw new SemanticException( String.format( "alias '%s' can only be matched with one time series", resultColumn.getAlias())); @@ -467,6 +469,12 @@ public class Analyzer { : null; alias = hasAlias ? resultColumn.getAlias() : alias; outputExpressions.add(new Pair<>(expressionWithoutAlias, alias)); + if (queryStatement.isGroupByLevel() + && resultColumn.getExpression() instanceof FunctionExpression) { + queryStatement + .getGroupByLevelComponent() + .updateIsCountStar((FunctionExpression) resultColumn.getExpression()); + } ExpressionAnalyzer.updateTypeProvider(expressionWithoutAlias, typeProvider); expressionWithoutAlias.inferTypes(typeProvider); paginationController.consumeLimit(); @@ -668,8 +676,11 @@ public class Analyzer { GroupByLevelController groupByLevelController = new GroupByLevelController( queryStatement.getGroupByLevelComponent().getLevels(), typeProvider); - for (Pair<Expression, String> measurementWithAlias : outputExpressions) { - groupByLevelController.control(measurementWithAlias.left, measurementWithAlias.right); + for (int i = 0; i < outputExpressions.size(); i++) { + Pair<Expression, String> measurementWithAlias = outputExpressions.get(i); + boolean isCountStar = queryStatement.getGroupByLevelComponent().isCountStar(i); + groupByLevelController.control( + isCountStar, measurementWithAlias.left, measurementWithAlias.right); } Map<Expression, Set<Expression>> rawGroupByLevelExpressions = groupByLevelController.getGroupedPathMap(); diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ExpressionAnalyzer.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ExpressionAnalyzer.java index 1cee37386b..0f0c2b2720 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ExpressionAnalyzer.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ExpressionAnalyzer.java @@ -289,6 +289,7 @@ public class ExpressionAnalyzer { // removing all wildcards. We use actualExpressions to collect them. List<List<Expression>> childExpressionsList = new ArrayList<>(); cartesianProduct(extendedExpressions, childExpressionsList, 0, new ArrayList<>()); + return reconstructFunctionExpressions((FunctionExpression) expression, childExpressionsList); } else if (expression instanceof TimeSeriesOperand) { PartialPath path = ((TimeSeriesOperand) expression).getPath(); diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/GroupByLevelController.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/GroupByLevelController.java index 31a1715e06..a583805e33 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/GroupByLevelController.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/GroupByLevelController.java @@ -74,14 +74,14 @@ public class GroupByLevelController { this.typeProvider = typeProvider; } - public void control(Expression expression, String alias) { + public void control(boolean isCountStar, Expression expression, String alias) { if (!(expression instanceof FunctionExpression && expression.isBuiltInAggregationFunctionExpression())) { throw new SemanticException(expression + " can't be used in group by level."); } PartialPath rawPath = ((TimeSeriesOperand) expression.getExpressions().get(0)).getPath(); - PartialPath groupedPath = generatePartialPathByLevel(rawPath.getNodes(), levels); + PartialPath groupedPath = generatePartialPathByLevel(isCountStar, rawPath.getNodes(), levels); checkDatatypeConsistency( groupedPath.getFullPath(), ((FunctionExpression) expression).getFunctionName(), rawPath); @@ -178,7 +178,8 @@ public class GroupByLevelController { * * @return result partial path */ - public PartialPath generatePartialPathByLevel(String[] nodes, int[] pathLevels) { + public PartialPath generatePartialPathByLevel( + boolean isCountStar, String[] nodes, int[] pathLevels) { Set<Integer> levelSet = new HashSet<>(); for (int level : pathLevels) { levelSet.add(level); @@ -194,7 +195,11 @@ public class GroupByLevelController { transformedNodes.add(IoTDBConstant.ONE_LEVEL_PATH_WILDCARD); } } - transformedNodes.add(nodes[nodes.length - 1]); + if (isCountStar) { + transformedNodes.add(IoTDBConstant.ONE_LEVEL_PATH_WILDCARD); + } else { + transformedNodes.add(nodes[nodes.length - 1]); + } return new PartialPath(transformedNodes.toArray(new String[0])); } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/component/GroupByLevelComponent.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/component/GroupByLevelComponent.java index fa7192a0fe..33db43633b 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/component/GroupByLevelComponent.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/component/GroupByLevelComponent.java @@ -19,12 +19,17 @@ package org.apache.iotdb.db.mpp.plan.statement.component; +import org.apache.iotdb.db.mpp.plan.expression.multi.FunctionExpression; import org.apache.iotdb.db.mpp.plan.statement.StatementNode; +import java.util.ArrayList; +import java.util.List; + /** This class maintains information of {@code GROUP BY LEVEL} clause. */ public class GroupByLevelComponent extends StatementNode { protected int[] levels; + protected List<Boolean> isCountStar = new ArrayList<>(); public int[] getLevels() { return levels; @@ -33,4 +38,12 @@ public class GroupByLevelComponent extends StatementNode { public void setLevels(int[] levels) { this.levels = levels; } + + public void updateIsCountStar(FunctionExpression rawExpression) { + isCountStar.add(rawExpression.isCountStar()); + } + + public boolean isCountStar(int i) { + return isCountStar.get(i); + } }
