[CALCITE-2012] Replace LocalInterval by Interval in Druid adapter
Project: http://git-wip-us.apache.org/repos/asf/calcite/repo Commit: http://git-wip-us.apache.org/repos/asf/calcite/commit/20ade9d2 Tree: http://git-wip-us.apache.org/repos/asf/calcite/tree/20ade9d2 Diff: http://git-wip-us.apache.org/repos/asf/calcite/diff/20ade9d2 Branch: refs/heads/master Commit: 20ade9d266d969f230be7b4e0062db17757e53a3 Parents: 830801b Author: Jesus Camacho Rodriguez <[email protected]> Authored: Fri Oct 13 10:29:27 2017 -0700 Committer: Jesus Camacho Rodriguez <[email protected]> Committed: Tue Nov 14 15:37:02 2017 -0800 ---------------------------------------------------------------------- .../adapter/druid/DruidConnectionImpl.java | 4 +- .../adapter/druid/DruidDateTimeUtils.java | 16 +- .../calcite/adapter/druid/DruidQuery.java | 14 +- .../calcite/adapter/druid/DruidRules.java | 3 +- .../calcite/adapter/druid/DruidTable.java | 17 +- .../adapter/druid/DruidTableFactory.java | 8 +- .../calcite/adapter/druid/LocalInterval.java | 98 ------- .../org/apache/calcite/test/DruidAdapterIT.java | 276 +++++++++---------- .../calcite/test/DruidDateRangeRulesTest.java | 32 +-- .../test/resources/druid-foodmart-model.json | 2 +- druid/src/test/resources/druid-wiki-model.json | 2 +- .../resources/druid-wiki-no-columns-model.json | 2 +- site/_docs/druid_adapter.md | 4 +- 13 files changed, 198 insertions(+), 280 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/calcite/blob/20ade9d2/druid/src/main/java/org/apache/calcite/adapter/druid/DruidConnectionImpl.java ---------------------------------------------------------------------- diff --git a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidConnectionImpl.java b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidConnectionImpl.java index 1951396..91fdf90 100644 --- a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidConnectionImpl.java +++ b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidConnectionImpl.java @@ -42,6 +42,8 @@ import com.google.common.base.Preconditions; import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; +import org.joda.time.Interval; + import java.io.ByteArrayInputStream; import java.io.IOException; import java.io.InputStream; @@ -493,7 +495,7 @@ class DruidConnectionImpl implements DruidConnection { /** Reads segment metadata, and populates a list of columns and metrics. */ void metadata(String dataSourceName, String timestampColumnName, - List<LocalInterval> intervals, + List<Interval> intervals, Map<String, SqlTypeName> fieldBuilder, Set<String> metricNameBuilder, Map<String, List<ComplexMetric>> complexMetrics) { final String url = this.url + "/druid/v2/?pretty"; http://git-wip-us.apache.org/repos/asf/calcite/blob/20ade9d2/druid/src/main/java/org/apache/calcite/adapter/druid/DruidDateTimeUtils.java ---------------------------------------------------------------------- diff --git a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidDateTimeUtils.java b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidDateTimeUtils.java index fb69353..562eb4a 100644 --- a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidDateTimeUtils.java +++ b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidDateTimeUtils.java @@ -38,6 +38,8 @@ import com.google.common.collect.Lists; import com.google.common.collect.Range; import com.google.common.collect.TreeRangeSet; +import org.joda.time.Interval; +import org.joda.time.chrono.ISOChronology; import org.slf4j.Logger; import java.util.ArrayList; @@ -56,11 +58,11 @@ public class DruidDateTimeUtils { } /** - * Generates a list of {@link LocalInterval}s equivalent to a given + * Generates a list of {@link Interval}s equivalent to a given * expression. Assumes that all the predicates in the input * reference a single column: the timestamp column. */ - public static List<LocalInterval> createInterval(RexNode e, String timeZone) { + public static List<Interval> createInterval(RexNode e, String timeZone) { final List<Range<TimestampString>> ranges = extractRanges(e, TimeZone.getTimeZone(timeZone), false); if (ranges == null) { @@ -78,11 +80,11 @@ public class DruidDateTimeUtils { ImmutableList.<Range>copyOf(condensedRanges.asRanges())); } - protected static List<LocalInterval> toInterval( + protected static List<Interval> toInterval( List<Range<TimestampString>> ranges) { - List<LocalInterval> intervals = Lists.transform(ranges, - new Function<Range<TimestampString>, LocalInterval>() { - public LocalInterval apply(Range<TimestampString> range) { + List<Interval> intervals = Lists.transform(ranges, + new Function<Range<TimestampString>, Interval>() { + public Interval apply(Range<TimestampString> range) { if (!range.hasLowerBound() && !range.hasUpperBound()) { return DruidTable.DEFAULT_INTERVAL; } @@ -100,7 +102,7 @@ public class DruidDateTimeUtils { && range.upperBoundType() == BoundType.CLOSED) { end++; } - return LocalInterval.create(start, end); + return new Interval(start, end, ISOChronology.getInstanceUTC()); } }); if (LOGGER.isDebugEnabled()) { http://git-wip-us.apache.org/repos/asf/calcite/blob/20ade9d2/druid/src/main/java/org/apache/calcite/adapter/druid/DruidQuery.java ---------------------------------------------------------------------- diff --git a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidQuery.java b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidQuery.java index 4cfee48..3a7f25a 100644 --- a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidQuery.java +++ b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidQuery.java @@ -73,6 +73,8 @@ import com.google.common.collect.ImmutableList; import com.google.common.collect.Iterables; import com.google.common.collect.Sets; +import org.joda.time.Interval; + import java.io.IOException; import java.io.StringWriter; import java.math.BigDecimal; @@ -95,7 +97,7 @@ public class DruidQuery extends AbstractRelNode implements BindableRel { final RelOptTable table; final DruidTable druidTable; - final ImmutableList<LocalInterval> intervals; + final ImmutableList<Interval> intervals; final ImmutableList<RelNode> rels; private static final Pattern VALID_SIG = Pattern.compile("sf?p?(a?|ao)l?"); @@ -115,7 +117,7 @@ public class DruidQuery extends AbstractRelNode implements BindableRel { */ protected DruidQuery(RelOptCluster cluster, RelTraitSet traitSet, RelOptTable table, DruidTable druidTable, - List<LocalInterval> intervals, List<RelNode> rels) { + List<Interval> intervals, List<RelNode> rels) { super(cluster, traitSet); this.table = table; this.druidTable = druidTable; @@ -292,7 +294,7 @@ public class DruidQuery extends AbstractRelNode implements BindableRel { /** Creates a DruidQuery. */ private static DruidQuery create(RelOptCluster cluster, RelTraitSet traitSet, - RelOptTable table, DruidTable druidTable, List<LocalInterval> intervals, + RelOptTable table, DruidTable druidTable, List<Interval> intervals, List<RelNode> rels) { return new DruidQuery(cluster, traitSet, table, druidTable, intervals, rels); } @@ -307,7 +309,7 @@ public class DruidQuery extends AbstractRelNode implements BindableRel { /** Extends a DruidQuery. */ public static DruidQuery extendQuery(DruidQuery query, - List<LocalInterval> intervals) { + List<Interval> intervals) { return DruidQuery.create(query.getCluster(), query.getTraitSet(), query.getTable(), query.druidTable, intervals, query.rels); } @@ -999,7 +1001,7 @@ public class DruidQuery extends AbstractRelNode implements BindableRel { if (o instanceof String) { String s = (String) o; generator.writeString(s); - } else if (o instanceof LocalInterval) { + } else if (o instanceof Interval) { generator.writeString(o.toString()); } else if (o instanceof Integer) { Integer i = (Integer) o; @@ -1015,7 +1017,7 @@ public class DruidQuery extends AbstractRelNode implements BindableRel { /** Generates a JSON string to query metadata about a data source. */ static String metadataQuery(String dataSourceName, - List<LocalInterval> intervals) { + List<Interval> intervals) { final StringWriter sw = new StringWriter(); final JsonFactory factory = new JsonFactory(); try { http://git-wip-us.apache.org/repos/asf/calcite/blob/20ade9d2/druid/src/main/java/org/apache/calcite/adapter/druid/DruidRules.java ---------------------------------------------------------------------- diff --git a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidRules.java b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidRules.java index baf5e41..7a22012 100644 --- a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidRules.java +++ b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidRules.java @@ -68,6 +68,7 @@ import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; import com.google.common.collect.Lists; +import org.joda.time.Interval; import org.slf4j.Logger; import java.util.ArrayList; @@ -233,7 +234,7 @@ public class DruidRules { return; } final List<RexNode> residualPreds = new ArrayList<>(triple.getRight()); - List<LocalInterval> intervals = null; + List<Interval> intervals = null; if (!triple.getLeft().isEmpty()) { intervals = DruidDateTimeUtils.createInterval( RexUtil.composeConjunction(rexBuilder, triple.getLeft(), false), http://git-wip-us.apache.org/repos/asf/calcite/blob/20ade9d2/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTable.java ---------------------------------------------------------------------- diff --git a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTable.java b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTable.java index 3789a36..10e0466 100644 --- a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTable.java +++ b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTable.java @@ -42,6 +42,10 @@ import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; +import org.joda.time.DateTime; +import org.joda.time.Interval; +import org.joda.time.chrono.ISOChronology; + import java.util.ArrayList; import java.util.List; import java.util.Map; @@ -53,14 +57,15 @@ import java.util.Set; public class DruidTable extends AbstractTable implements TranslatableTable { public static final String DEFAULT_TIMESTAMP_COLUMN = "__time"; - public static final LocalInterval DEFAULT_INTERVAL = - LocalInterval.create("1900-01-01", "3000-01-01"); + public static final Interval DEFAULT_INTERVAL = + new Interval(new DateTime("1900-01-01", ISOChronology.getInstanceUTC()), + new DateTime("3000-01-01", ISOChronology.getInstanceUTC())); final DruidSchema schema; final String dataSource; final RelProtoDataType protoRowType; final ImmutableSet<String> metricFieldNames; - final ImmutableList<LocalInterval> intervals; + final ImmutableList<Interval> intervals; final String timestampFieldName; final ImmutableMap<String, List<ComplexMetric>> complexMetrics; final ImmutableMap<String, SqlTypeName> allFields; @@ -77,7 +82,7 @@ public class DruidTable extends AbstractTable implements TranslatableTable { */ public DruidTable(DruidSchema schema, String dataSource, RelProtoDataType protoRowType, Set<String> metricFieldNames, - String timestampFieldName, List<LocalInterval> intervals, + String timestampFieldName, List<Interval> intervals, Map<String, List<ComplexMetric>> complexMetrics, Map<String, SqlTypeName> allFields) { this.timestampFieldName = Preconditions.checkNotNull(timestampFieldName); this.schema = Preconditions.checkNotNull(schema); @@ -107,7 +112,7 @@ public class DruidTable extends AbstractTable implements TranslatableTable { * @return A table */ static Table create(DruidSchema druidSchema, String dataSourceName, - List<LocalInterval> intervals, Map<String, SqlTypeName> fieldMap, + List<Interval> intervals, Map<String, SqlTypeName> fieldMap, Set<String> metricNameSet, String timestampColumnName, DruidConnectionImpl connection, Map<String, List<ComplexMetric>> complexMetrics) { assert connection != null; @@ -132,7 +137,7 @@ public class DruidTable extends AbstractTable implements TranslatableTable { * @return A table */ static Table create(DruidSchema druidSchema, String dataSourceName, - List<LocalInterval> intervals, Map<String, SqlTypeName> fieldMap, + List<Interval> intervals, Map<String, SqlTypeName> fieldMap, Set<String> metricNameSet, String timestampColumnName, Map<String, List<ComplexMetric>> complexMetrics) { final ImmutableMap<String, SqlTypeName> fields = http://git-wip-us.apache.org/repos/asf/calcite/blob/20ade9d2/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTableFactory.java ---------------------------------------------------------------------- diff --git a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTableFactory.java b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTableFactory.java index d34e000..add0136 100644 --- a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTableFactory.java +++ b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidTableFactory.java @@ -25,6 +25,9 @@ import org.apache.calcite.util.Util; import com.google.common.collect.ImmutableList; +import org.joda.time.Interval; +import org.joda.time.chrono.ISOChronology; + import java.util.ArrayList; import java.util.HashMap; import java.util.LinkedHashMap; @@ -121,9 +124,10 @@ public class DruidTableFactory implements TableFactory { } } final Object interval = operand.get("interval"); - final List<LocalInterval> intervals; + final List<Interval> intervals; if (interval instanceof String) { - intervals = ImmutableList.of(LocalInterval.create((String) interval)); + intervals = ImmutableList.of( + new Interval((String) interval, ISOChronology.getInstanceUTC())); } else { intervals = null; } http://git-wip-us.apache.org/repos/asf/calcite/blob/20ade9d2/druid/src/main/java/org/apache/calcite/adapter/druid/LocalInterval.java ---------------------------------------------------------------------- diff --git a/druid/src/main/java/org/apache/calcite/adapter/druid/LocalInterval.java b/druid/src/main/java/org/apache/calcite/adapter/druid/LocalInterval.java deleted file mode 100644 index 2485281..0000000 --- a/druid/src/main/java/org/apache/calcite/adapter/druid/LocalInterval.java +++ /dev/null @@ -1,98 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to you under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.calcite.adapter.druid; - -import org.joda.time.DateTime; -import org.joda.time.Instant; -import org.joda.time.Interval; -import org.joda.time.LocalDateTime; -import org.joda.time.chrono.ISOChronology; - -/** - * Similar to {@link Interval} but the end points are {@link LocalDateTime} - * not {@link Instant}. - */ -public class LocalInterval { - private final long start; - private final long end; - - /** Creates a LocalInterval. */ - private LocalInterval(long start, long end) { - this.start = start; - this.end = end; - } - - /** Creates a LocalInterval based on two {@link DateTime} values. */ - public static LocalInterval create(DateTime start, DateTime end) { - return new LocalInterval(start.getMillis(), end.getMillis()); - } - - /** Creates a LocalInterval based on millisecond start end end points. */ - public static LocalInterval create(long start, long end) { - return new LocalInterval(start, end); - } - - /** Creates a LocalInterval based on an interval string. */ - public static LocalInterval create(String intervalString) { - Interval i = new Interval(intervalString, ISOChronology.getInstanceUTC()); - return new LocalInterval(i.getStartMillis(), i.getEndMillis()); - } - - /** Creates a LocalInterval based on start and end time strings. */ - public static LocalInterval create(String start, String end) { - return create( - new DateTime(start, ISOChronology.getInstanceUTC()), - new DateTime(end, ISOChronology.getInstanceUTC())); - } - - /** Writes a value such as "1900-01-01T00:00:00.000/2015-10-12T00:00:00.000". - * Note that there are no "Z"s; the value is in the (unspecified) local - * time zone, not UTC. */ - @Override public String toString() { - final LocalDateTime start = - new LocalDateTime(this.start, ISOChronology.getInstanceUTC()); - final LocalDateTime end = - new LocalDateTime(this.end, ISOChronology.getInstanceUTC()); - return start + "/" + end; - } - - @Override public int hashCode() { - int result = 97; - result = 31 * result + ((int) (start ^ (start >>> 32))); - result = 31 * result + ((int) (end ^ (end >>> 32))); - return result; - } - - @Override public boolean equals(Object o) { - return o == this - || o instanceof LocalInterval - && start == ((LocalInterval) o).start - && end == ((LocalInterval) o).end; - } - - /** Analogous to {@link Interval#getStartMillis}. */ - public long getStartMillis() { - return start; - } - - /** Analogous to {@link Interval#getEndMillis()}. */ - public long getEndMillis() { - return end; - } -} - -// End LocalInterval.java
