Repository: calcite Updated Branches: refs/heads/master f4b4ee115 -> 884f01066
Add Orinoco schema (streaming retail data), accessible from Quidem scripts Project: http://git-wip-us.apache.org/repos/asf/calcite/repo Commit: http://git-wip-us.apache.org/repos/asf/calcite/commit/cbbaddf5 Tree: http://git-wip-us.apache.org/repos/asf/calcite/tree/cbbaddf5 Diff: http://git-wip-us.apache.org/repos/asf/calcite/diff/cbbaddf5 Branch: refs/heads/master Commit: cbbaddf5449d7547e7550d08caaadf36537cd7fd Parents: f4b4ee1 Author: Julian Hyde <[email protected]> Authored: Fri Feb 19 13:38:00 2016 -0800 Committer: Julian Hyde <[email protected]> Committed: Sun Feb 21 01:42:28 2016 -0800 ---------------------------------------------------------------------- .../org/apache/calcite/test/CalciteAssert.java | 9 +++- .../java/org/apache/calcite/test/JdbcTest.java | 6 +++ .../org/apache/calcite/test/StreamTest.java | 56 ++++++++++++++++++-- core/src/test/resources/sql/agg.iq | 31 +++++++++++ 4 files changed, 96 insertions(+), 6 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/calcite/blob/cbbaddf5/core/src/test/java/org/apache/calcite/test/CalciteAssert.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/calcite/test/CalciteAssert.java b/core/src/test/java/org/apache/calcite/test/CalciteAssert.java index 03db1b9..e15664f 100644 --- a/core/src/test/java/org/apache/calcite/test/CalciteAssert.java +++ b/core/src/test/java/org/apache/calcite/test/CalciteAssert.java @@ -678,6 +678,12 @@ public class CalciteAssert { case LINGUAL: return rootSchema.add("SALES", new ReflectiveSchema(new JdbcTest.LingualSchema())); + case ORINOCO: + final SchemaPlus orinoco = rootSchema.add("ORINOCO", new AbstractSchema()); + orinoco.add("ORDERS", + new StreamTest.OrdersHistoryTable( + StreamTest.OrdersStreamTableFactory.getRowList())); + return orinoco; case POST: final SchemaPlus post = rootSchema.add("POST", new AbstractSchema()); post.add("EMP", @@ -1581,7 +1587,8 @@ public class CalciteAssert { JDBC_SCOTT, SCOTT, LINGUAL, - POST + POST, + ORINOCO } /** Converts a {@link ResultSet} to string. */ http://git-wip-us.apache.org/repos/asf/calcite/blob/cbbaddf5/core/src/test/java/org/apache/calcite/test/JdbcTest.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/calcite/test/JdbcTest.java b/core/src/test/java/org/apache/calcite/test/JdbcTest.java index fc4ca15..d7b635c 100644 --- a/core/src/test/java/org/apache/calcite/test/JdbcTest.java +++ b/core/src/test/java/org/apache/calcite/test/JdbcTest.java @@ -4846,6 +4846,12 @@ public class JdbcTest { new ReflectiveSchemaTest.CatchallSchema())) .connect(); } + if (name.equals("orinoco")) { + return CalciteAssert.that() + .with(CalciteAssert.SchemaSpec.ORINOCO) + .withDefaultSchema("ORINOCO") + .connect(); + } if (name.equals("seq")) { final Connection connection = CalciteAssert.that() .withSchema("s", new AbstractSchema()) http://git-wip-us.apache.org/repos/asf/calcite/blob/cbbaddf5/core/src/test/java/org/apache/calcite/test/StreamTest.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/calcite/test/StreamTest.java b/core/src/test/java/org/apache/calcite/test/StreamTest.java index 649db3d..99403bd 100644 --- a/core/src/test/java/org/apache/calcite/test/StreamTest.java +++ b/core/src/test/java/org/apache/calcite/test/StreamTest.java @@ -261,6 +261,33 @@ public class StreamTest { "ROWTIME=2015-02-15 10:24:45; ORDERID=3; SUPPLIERID=1")); } + @Ignore + @Test public void testTumbleViaOver() { + String sql = "WITH HourlyOrderTotals (rowtime, productId, c, su) AS (\n" + + " SELECT FLOOR(rowtime TO HOUR),\n" + + " productId,\n" + + " COUNT(*),\n" + + " SUM(units)\n" + + " FROM Orders\n" + + " GROUP BY FLOOR(rowtime TO HOUR), productId)\n" + + "SELECT STREAM rowtime,\n" + + " productId,\n" + + " SUM(su) OVER w AS su,\n" + + " SUM(c) OVER w AS c\n" + + "FROM HourlyTotals\n" + + "WINDOW w AS (\n" + + " ORDER BY rowtime\n" + + " PARTITION BY productId\n" + + " RANGE INTERVAL '2' HOUR PRECEDING)\n"; + String sql2 = "" + + "SELECT STREAM rowtime, productId, SUM(units) AS su, COUNT(*) AS c\n" + + "FROM Orders\n" + + "GROUP BY TUMBLE(rowtime, INTERVAL '1' HOUR)"; + // sql and sql2 should give same result + CalciteAssert.model(STREAM_JOINS_MODEL) + .query(sql); + } + private Function<ResultSet, Void> startsWith(String... rows) { final ImmutableList<String> rowList = ImmutableList.copyOf(rows); return new Function<ResultSet, Void>() { @@ -289,7 +316,7 @@ public class StreamTest { * Base table for the Orders table. Manages the base schema used for the test tables and common * functions. */ - private abstract static class BaseOrderStreamTable implements ScannableTable, StreamableTable { + private abstract static class BaseOrderStreamTable implements ScannableTable { protected final RelProtoDataType protoRowType = new RelProtoDataType() { public RelDataType apply(RelDataTypeFactory a0) { return a0.builder() @@ -325,6 +352,10 @@ public class StreamTest { public Table create(SchemaPlus schema, String name, Map<String, Object> operand, RelDataType rowType) { + return new OrdersTable(getRowList()); + } + + public static ImmutableList<Object[]> getRowList() { final Object[][] rows = { {ts(10, 15, 0), 1, "paint", 10}, {ts(10, 24, 15), 2, "paper", 5}, @@ -332,16 +363,17 @@ public class StreamTest { {ts(10, 58, 0), 4, "paint", 3}, {ts(11, 10, 0), 5, "paint", 3} }; - return new OrdersTable(ImmutableList.copyOf(rows)); + return ImmutableList.copyOf(rows); } - private Object ts(int h, int m, int s) { + private static Object ts(int h, int m, int s) { return DateTimeUtils.unixTimestamp(2015, 2, 15, h, m, s); } } /** Table representing the ORDERS stream. */ - public static class OrdersTable extends BaseOrderStreamTable { + public static class OrdersTable extends BaseOrderStreamTable + implements StreamableTable { private final ImmutableList<Object[]> rows; public OrdersTable(ImmutableList<Object[]> rows) { @@ -386,7 +418,8 @@ public class StreamTest { /** * Table representing an infinitely larger ORDERS stream. */ - public static class InfiniteOrdersTable extends BaseOrderStreamTable { + public static class InfiniteOrdersTable extends BaseOrderStreamTable + implements StreamableTable { public Enumerable<Object[]> scan(DataContext root) { return Linq4j.asEnumerable(new Iterable<Object[]>() { @Override public Iterator<Object[]> iterator() { @@ -412,6 +445,19 @@ public class StreamTest { } } + /** Table representing the history of the ORDERS stream. */ + public static class OrdersHistoryTable extends BaseOrderStreamTable { + private final ImmutableList<Object[]> rows; + + public OrdersHistoryTable(ImmutableList<Object[]> rows) { + this.rows = rows; + } + + public Enumerable<Object[]> scan(DataContext root) { + return Linq4j.asEnumerable(rows); + } + } + /** * Mocks a simple relation to use for stream joining test. */ http://git-wip-us.apache.org/repos/asf/calcite/blob/cbbaddf5/core/src/test/resources/sql/agg.iq ---------------------------------------------------------------------- diff --git a/core/src/test/resources/sql/agg.iq b/core/src/test/resources/sql/agg.iq index 3b1f58e..9390b02 100644 --- a/core/src/test/resources/sql/agg.iq +++ b/core/src/test/resources/sql/agg.iq @@ -1493,6 +1493,37 @@ EnumerableCalc(expr#0..2=[{inputs}], JOB=[$t0], SUM_SAL=[$t2], DEPTNO=[$t1]) !plan !} +!use orinoco + +# FLOOR to achieve a 2-hour window +select floor(rowtime to hour) as rowtime, count(*) as c +from Orders +group by floor(rowtime to hour); ++---------------------+---+ +| ROWTIME | C | ++---------------------+---+ +| 2015-02-15 10:00:00 | 4 | +| 2015-02-15 11:00:00 | 1 | ++---------------------+---+ +(2 rows) + +!ok + +# FLOOR applied to intervals, to achieve a 2-hour window +select rowtime, count(*) as c +from ( + select timestamp '1970-1-1 0:0:0' + (floor(timestamp '1970-1-1 0:0:0' + ((rowtime - timestamp '1970-1-1 0:0:0') second) / 2 to hour) - timestamp '1970-1-1 0:0:0') second * 2 as rowtime + from Orders) +group by rowtime; ++---------------------+---+ +| ROWTIME | C | ++---------------------+---+ +| 2015-02-15 10:00:00 | 5 | ++---------------------+---+ +(1 row) + +!ok + # [CALCITE-729] IndexOutOfBoundsException in ROLLUP query on JDBC data source !use jdbc_scott select deptno, job, count(*) as c
