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

Reply via email to