Repository: calcite Updated Branches: refs/heads/master 61f1cf925 -> f6825f079
[CALCITE-1959] Reduce the amount of metadata and tableName calls in Druid (Zain Humayun) Close apache/calcite#524 Project: http://git-wip-us.apache.org/repos/asf/calcite/repo Commit: http://git-wip-us.apache.org/repos/asf/calcite/commit/f6825f07 Tree: http://git-wip-us.apache.org/repos/asf/calcite/tree/f6825f07 Diff: http://git-wip-us.apache.org/repos/asf/calcite/diff/f6825f07 Branch: refs/heads/master Commit: f6825f079fe6812d93fd4a90f669c8c33af0a580 Parents: 61f1cf9 Author: Zain Humayun <[email protected]> Authored: Mon Aug 21 13:39:30 2017 -0700 Committer: Jesus Camacho Rodriguez <[email protected]> Committed: Wed Aug 23 10:36:29 2017 -0700 ---------------------------------------------------------------------- .../calcite/adapter/druid/DruidSchema.java | 43 ++++++++++++-------- .../calcite/adapter/druid/DruidTable.java | 39 ++++++++++++++---- .../adapter/druid/DruidTableFactory.java | 22 +++++----- .../org/apache/calcite/test/DruidAdapterIT.java | 10 +++++ 4 files changed, 78 insertions(+), 36 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/calcite/blob/f6825f07/druid/src/main/java/org/apache/calcite/adapter/druid/DruidSchema.java ---------------------------------------------------------------------- diff --git a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidSchema.java b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidSchema.java index bc344f3..87717f2 100644 --- a/druid/src/main/java/org/apache/calcite/adapter/druid/DruidSchema.java +++ b/druid/src/main/java/org/apache/calcite/adapter/druid/DruidSchema.java @@ -42,6 +42,7 @@ public class DruidSchema extends AbstractSchema { final String url; final String coordinatorUrl; private final boolean discoverTables; + private Map<String, Table> tableMap = null; /** * Creates a Druid schema. @@ -63,23 +64,31 @@ public class DruidSchema extends AbstractSchema { if (!discoverTables) { return ImmutableMap.of(); } - final DruidConnectionImpl connection = - new DruidConnectionImpl(url, coordinatorUrl); - return Compatible.INSTANCE.asMap( - ImmutableSet.copyOf(connection.tableNames()), - CacheBuilder.newBuilder() - .build(new CacheLoader<String, Table>() { - public Table load(@Nonnull String tableName) throws Exception { - final Map<String, SqlTypeName> fieldMap = new LinkedHashMap<>(); - final Set<String> metricNameSet = new LinkedHashSet<>(); - final Map<String, List<ComplexMetric>> complexMetrics = new HashMap<>(); - connection.metadata(tableName, DruidTable.DEFAULT_TIMESTAMP_COLUMN, - null, fieldMap, metricNameSet, complexMetrics); - return DruidTable.create(DruidSchema.this, tableName, null, - fieldMap, metricNameSet, DruidTable.DEFAULT_TIMESTAMP_COLUMN, connection, - complexMetrics); - } - })); + + if (tableMap == null) { + final DruidConnectionImpl connection = new DruidConnectionImpl(url, coordinatorUrl); + Set<String> tableNames = connection.tableNames(); + + tableMap = Compatible.INSTANCE.asMap( + ImmutableSet.copyOf(tableNames), + CacheBuilder.newBuilder() + .build(new CacheLoader<String, Table>() { + @Override public Table load(@Nonnull String tableName) throws Exception { + final Map<String, SqlTypeName> fieldMap = new LinkedHashMap<>(); + final Set<String> metricNameSet = new LinkedHashSet<>(); + final Map<String, List<ComplexMetric>> complexMetrics = new HashMap<>(); + + connection.metadata(tableName, DruidTable.DEFAULT_TIMESTAMP_COLUMN, + null, fieldMap, metricNameSet, complexMetrics); + + return DruidTable.create(DruidSchema.this, tableName, null, + fieldMap, metricNameSet, DruidTable.DEFAULT_TIMESTAMP_COLUMN, + connection, complexMetrics); + } + })); + } + + return tableMap; } } http://git-wip-us.apache.org/repos/asf/calcite/blob/f6825f07/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 7597773..8cff818 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 @@ -101,23 +101,44 @@ public class DruidTable extends AbstractTable implements TranslatableTable { * @param metricNameSet Mutable set of metric names; * may be partially populated already * @param timestampColumnName Name of timestamp column, or null - * @param connection If not null, use this connection to find column - * definitions + * @param connection connection used to find column definitions. Must be non-null. + * * @return A table */ static Table create(DruidSchema druidSchema, String dataSourceName, List<LocalInterval> intervals, Map<String, SqlTypeName> fieldMap, Set<String> metricNameSet, String timestampColumnName, DruidConnectionImpl connection, Map<String, List<ComplexMetric>> complexMetrics) { - if (connection != null) { - connection.metadata(dataSourceName, timestampColumnName, intervals, - fieldMap, metricNameSet, complexMetrics); - } + assert connection != null; + + connection.metadata(dataSourceName, timestampColumnName, intervals, + fieldMap, metricNameSet, complexMetrics); + + return DruidTable.create(druidSchema, dataSourceName, intervals, fieldMap, + metricNameSet, timestampColumnName, complexMetrics); + } + + /** Creates a {@link DruidTable} + * + * @param druidSchema Druid schema + * @param dataSourceName Data source name in Druid, also table name + * @param intervals Intervals, or null to use default + * @param fieldMap Mutable map of fields (dimensions plus metrics); + * may be partially populated already + * @param metricNameSet Mutable set of metric names; + * may be partially populated already + * @param timestampColumnName Name of timestamp column, or null + * @return A table + */ + static Table create(DruidSchema druidSchema, String dataSourceName, + List<LocalInterval> intervals, Map<String, SqlTypeName> fieldMap, + Set<String> metricNameSet, String timestampColumnName, + Map<String, List<ComplexMetric>> complexMetrics) { final ImmutableMap<String, SqlTypeName> fields = - ImmutableMap.copyOf(fieldMap); + ImmutableMap.copyOf(fieldMap); return new DruidTable(druidSchema, dataSourceName, - new MapRelProtoDataType(fields), ImmutableSet.copyOf(metricNameSet), - timestampColumnName, intervals, complexMetrics, fieldMap); + new MapRelProtoDataType(fields), ImmutableSet.copyOf(metricNameSet), + timestampColumnName, intervals, complexMetrics, fieldMap); } /** http://git-wip-us.apache.org/repos/asf/calcite/blob/f6825f07/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 c83348e..d636ce8 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 @@ -120,13 +120,6 @@ public class DruidTableFactory implements TableFactory { } } } - final String dataSourceName = Util.first(dataSource, name); - DruidConnectionImpl c; - if (dimensionsRaw == null || metricsRaw == null) { - c = new DruidConnectionImpl(druidSchema.url, druidSchema.url.replace(":8082", ":8081")); - } else { - c = null; - } final Object interval = operand.get("interval"); final List<LocalInterval> intervals; if (interval instanceof String) { @@ -134,10 +127,19 @@ public class DruidTableFactory implements TableFactory { } else { intervals = null; } - return DruidTable.create(druidSchema, dataSourceName, intervals, - fieldBuilder, metricNameBuilder, timestampColumnName, c, complexMetrics); - } + final String dataSourceName = Util.first(dataSource, name); + + if (dimensionsRaw == null || metricsRaw == null) { + DruidConnectionImpl connection = new DruidConnectionImpl(druidSchema.url, + druidSchema.url.replace(":8082", ":8081")); + return DruidTable.create(druidSchema, dataSourceName, intervals, fieldBuilder, + metricNameBuilder, timestampColumnName, connection, complexMetrics); + } else { + return DruidTable.create(druidSchema, dataSourceName, intervals, fieldBuilder, + metricNameBuilder, timestampColumnName, complexMetrics); + } + } } // End DruidTableFactory.java http://git-wip-us.apache.org/repos/asf/calcite/blob/f6825f07/druid/src/test/java/org/apache/calcite/test/DruidAdapterIT.java ---------------------------------------------------------------------- diff --git a/druid/src/test/java/org/apache/calcite/test/DruidAdapterIT.java b/druid/src/test/java/org/apache/calcite/test/DruidAdapterIT.java index 0eed641..fd64309 100644 --- a/druid/src/test/java/org/apache/calcite/test/DruidAdapterIT.java +++ b/druid/src/test/java/org/apache/calcite/test/DruidAdapterIT.java @@ -17,11 +17,13 @@ package org.apache.calcite.test; import org.apache.calcite.adapter.druid.DruidQuery; +import org.apache.calcite.adapter.druid.DruidSchema; import org.apache.calcite.config.CalciteConnectionConfig; import org.apache.calcite.config.CalciteConnectionProperty; import org.apache.calcite.rel.RelNode; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rex.RexNode; +import org.apache.calcite.schema.impl.AbstractSchema; import org.apache.calcite.sql.fun.SqlStdOperatorTable; import org.apache.calcite.sql.type.SqlTypeName; import org.apache.calcite.tools.RelBuilder; @@ -3183,6 +3185,14 @@ public class DruidAdapterIT { + ":true},'context':{'druid.query.fetch':true}}")); } + /** + * Test to make sure that the mapping from a Table name to a Table returned from + * {@link org.apache.calcite.adapter.druid.DruidSchema} is always the same Java object. + * */ + @Test public void testTableMapReused() { + AbstractSchema schema = new DruidSchema("http://localhost:8082", "http://localhost:8081", true); + assert schema.getTable("wikiticker") == schema.getTable("wikiticker"); + } } // End DruidAdapterIT.java
