From d379529ac32a99736e06f2d5dfad0de9cdacf775 Mon Sep 17 00:00:00 2001 From: SURYA SUMANTH N Date: Thu, 20 May 2021 19:22:45 +0530 Subject: [PATCH] Query Concurrency Code Optimizations --- .../carbondata/CarbondataMetadataFactory.java | 2 +- .../plugin/datacenter/DataCenterMetadata.java | 2 +- .../prestosql/plugin/jdbc/JdbcMetadata.java | 2 +- .../hive/BackgroundHiveSplitLoader.java | 18 ++-- .../prestosql/plugin/hive/HiveMetadata.java | 43 ++++++--- .../plugin/hive/HiveMetadataFactory.java | 9 +- .../plugin/hive/HivePartitionManager.java | 21 +++-- .../plugin/hive/HiveSplitManager.java | 2 +- .../plugin/hive/HiveSplitSource.java | 2 +- .../io/prestosql/plugin/hive/HiveUtil.java | 6 +- .../hive/metastore/CachingHiveMetastore.java | 11 ++- .../SemiTransactionalHiveMetastore.java | 13 ++- .../statistics/HiveStatisticsProvider.java | 7 +- .../MetastoreHiveStatisticsProvider.java | 91 ++++++++++++++++--- .../statistics/TableColumnStatistics.java | 30 ++++++ .../plugin/hive/AbstractTestHive.java | 9 +- .../hive/TestBackgroundHiveSplitLoader.java | 18 +++- .../TestMetastoreHiveStatisticsProvider.java | 31 ++++--- .../io/prestosql/cost/TableScanStatsRule.java | 2 +- .../dispatcher/LocalDispatchQueryFactory.java | 5 +- .../java/io/prestosql/event/QueryMonitor.java | 26 +++++- .../execution/QueryStateMachine.java | 32 +++++++ .../prestosql/execution/QueryStateTimer.java | 35 +++++++ .../io/prestosql/execution/QueryStats.java | 18 ++++ .../execution/SqlQueryExecution.java | 2 + .../InternalResourceGroupManager.java | 12 +++ .../resourcegroups/ResourceGroupManager.java | 10 ++ .../scheduler/SqlQueryScheduler.java | 4 +- .../java/io/prestosql/metadata/Metadata.java | 4 +- .../prestosql/metadata/MetadataManager.java | 5 +- .../query/CachedSqlQueryExecution.java | 33 ++++++- .../queryeditorui/execution/Execution.java | 2 + .../io/prestosql/server/QueryResource.java | 2 + .../rule/RemoveUnsupportedDynamicFilters.java | 4 +- .../planner/iterative/rule/TablePushdown.java | 2 +- .../optimizations/AddReuseExchange.java | 2 +- .../sql/rewrite/ShowStatsRewrite.java | 2 +- .../InMemoryTransactionManager.java | 26 +++--- .../prestosql/execution/TestQueryStats.java | 2 + .../metadata/AbstractMockMetadata.java | 2 +- .../prestosql/server/TestBasicQueryInfo.java | 2 + .../prestosql/server/TestQueryStateInfo.java | 2 + .../single_query_info_response.json | 2 + .../connector/CachedConnectorMetadata.java | 8 +- .../spi/connector/ConnectorMetadata.java | 2 +- .../ClassLoaderSafeConnectorMetadata.java | 4 +- .../spi/statistics/TableStatistics.java | 33 ++++++- .../prestosql/plugin/tpcds/TpcdsMetadata.java | 2 +- .../tpcds/TestTpcdsMetadataStatistics.java | 8 +- .../prestosql/plugin/tpch/TpchMetadata.java | 2 +- .../plugin/tpch/TestTpchMetadata.java | 6 +- 51 files changed, 481 insertions(+), 139 deletions(-) create mode 100644 presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/TableColumnStatistics.java diff --git a/hetu-carbondata/src/main/java/io/hetu/core/plugin/carbondata/CarbondataMetadataFactory.java b/hetu-carbondata/src/main/java/io/hetu/core/plugin/carbondata/CarbondataMetadataFactory.java index 7cf7c3138..3abe8b88c 100755 --- a/hetu-carbondata/src/main/java/io/hetu/core/plugin/carbondata/CarbondataMetadataFactory.java +++ b/hetu-carbondata/src/main/java/io/hetu/core/plugin/carbondata/CarbondataMetadataFactory.java @@ -227,7 +227,7 @@ public class CarbondataMetadataFactory this.segmentInfoCodec, this.typeTranslator, this.hetuVersion, - new MetastoreHiveStatisticsProvider(metastore), + new MetastoreHiveStatisticsProvider(metastore, statsCache, samplePartitionCache), this.accessControlMetadataFactory.create(metastore), carbondataTableReader, this.carbondataTableStore, diff --git a/hetu-datacenter/src/main/java/io/hetu/core/plugin/datacenter/DataCenterMetadata.java b/hetu-datacenter/src/main/java/io/hetu/core/plugin/datacenter/DataCenterMetadata.java index b544617f6..91c9f6d60 100644 --- a/hetu-datacenter/src/main/java/io/hetu/core/plugin/datacenter/DataCenterMetadata.java +++ b/hetu-datacenter/src/main/java/io/hetu/core/plugin/datacenter/DataCenterMetadata.java @@ -242,7 +242,7 @@ public class DataCenterMetadata } @Override - public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint) + public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { Map columnHandles = getColumnHandles(session, tableHandle); String tableFullName = tableHandle.getSchemaPrefixedTableName(); diff --git a/presto-base-jdbc/src/main/java/io/prestosql/plugin/jdbc/JdbcMetadata.java b/presto-base-jdbc/src/main/java/io/prestosql/plugin/jdbc/JdbcMetadata.java index cf62c9ec6..993088f79 100644 --- a/presto-base-jdbc/src/main/java/io/prestosql/plugin/jdbc/JdbcMetadata.java +++ b/presto-base-jdbc/src/main/java/io/prestosql/plugin/jdbc/JdbcMetadata.java @@ -304,7 +304,7 @@ public class JdbcMetadata } @Override - public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint) + public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { JdbcTableHandle handle = (JdbcTableHandle) tableHandle; return jdbcClient.getTableStatistics(session, handle, constraint.getSummary()); diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/BackgroundHiveSplitLoader.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/BackgroundHiveSplitLoader.java index 85d6f851d..832ced586 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/BackgroundHiveSplitLoader.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/BackgroundHiveSplitLoader.java @@ -133,6 +133,7 @@ public class BackgroundHiveSplitLoader private final Deque> fileIterators = new ConcurrentLinkedDeque<>(); private final Optional validWriteIds; private final Supplier> dynamicFilterSupplier; + private final Configuration configuration; // Purpose of this lock: // * Write lock: when you need a consistent view across partitions, fileIterators, and hiveSplitSource. @@ -156,6 +157,7 @@ public class BackgroundHiveSplitLoader private Optional queryType; private Map queryInfo; private TypeManager typeManager; + private JobConf jobConf; private final Map cachedDynamicFilters = new ConcurrentHashMap<>(); @@ -194,6 +196,9 @@ public class BackgroundHiveSplitLoader this.queryType = requireNonNull(queryType, "queryType is null"); this.queryInfo = requireNonNull(queryInfo, "queryproperties is null"); this.partitions = new ConcurrentLazyQueue<>(getPrunedPartitions(partitions)); + Path path = new Path(getPartitionLocation(table, getPrunedPartitions(partitions).iterator().next().getPartition())); + configuration = hdfsEnvironment.getConfiguration(hdfsContext, path); + jobConf = ConfigurationUtils.toJobConf(configuration); } /** @@ -353,8 +358,7 @@ public class BackgroundHiveSplitLoader } Path path = new Path(getPartitionLocation(table, partition.getPartition())); - Configuration configuration = hdfsEnvironment.getConfiguration(hdfsContext, path); - InputFormat inputFormat = getInputFormat(configuration, schema, false); + InputFormat inputFormat = getInputFormat(configuration, schema, false, jobConf); FileSystem fs = hdfsEnvironment.getFileSystem(hdfsContext, path); boolean s3SelectPushdownEnabled = shouldEnablePushdownForTable(session, table, path.toString(), partition.getPartition()); @@ -371,11 +375,10 @@ public class BackgroundHiveSplitLoader // the splits must be generated using the file system for the target path // get the configuration for the target path -- it may be a different hdfs instance FileSystem targetFilesystem = hdfsEnvironment.getFileSystem(hdfsContext, targetPath); - JobConf targetJob = ConfigurationUtils.toJobConf(targetFilesystem.getConf()); - targetJob.setInputFormat(TextInputFormat.class); - targetInputFormat.configure(targetJob); - FileInputFormat.setInputPaths(targetJob, targetPath); - InputSplit[] targetSplits = targetInputFormat.getSplits(targetJob, 0); + jobConf.setInputFormat(TextInputFormat.class); + targetInputFormat.configure(jobConf); + FileInputFormat.setInputPaths(jobConf, targetPath); + InputSplit[] targetSplits = targetInputFormat.getSplits(jobConf, 0); InternalHiveSplitFactory splitFactory = new InternalHiveSplitFactory( targetFilesystem, @@ -437,7 +440,6 @@ public class BackgroundHiveSplitLoader throw new PrestoException(NOT_SUPPORTED, "Hive transactional tables in an input format with UseFileSplitsFromInputFormat annotation are not supported: " + inputFormat.getClass().getSimpleName()); } - JobConf jobConf = ConfigurationUtils.toJobConf(configuration); FileInputFormat.setInputPaths(jobConf, path); InputSplit[] splits = inputFormat.getSplits(jobConf, 0); diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveMetadata.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveMetadata.java index 66121749f..a6b472eb8 100755 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveMetadata.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveMetadata.java @@ -405,6 +405,11 @@ public class HiveMetadata return Optional.empty(); } + SchemaTableName schemaTableName = sourceTableHandle.getSchemaTableName(); + + Table table = metastore.getTable(new HiveIdentity(session), schemaTableName.getSchemaName(), schemaTableName.getTableName()) + .orElseThrow(() -> new TableNotFoundException(schemaTableName)); + List partitionColumns = sourceTableHandle.getPartitionColumns(); if (partitionColumns.isEmpty()) { return Optional.empty(); @@ -435,7 +440,7 @@ public class HiveMetadata Predicate> targetPredicate = convertToPredicate(targetTupleDomain); Constraint targetConstraint = new Constraint(targetTupleDomain, targetPredicate); Iterable> records = () -> - stream(partitionManager.getPartitions(metastore, new HiveIdentity(session), sourceTableHandle, targetConstraint).getPartitions()) + stream(partitionManager.getPartitions(metastore, new HiveIdentity(session), sourceTableHandle, targetConstraint, table).getPartitions()) .map(hivePartition -> IntStream.range(0, partitionColumns.size()) .mapToObj(fieldIdToColumnHandle::get) @@ -647,6 +652,12 @@ public class HiveMetadata .collect(toImmutableMap(HiveColumnHandle::getName, identity())); } + private Map getColumnHandles(Table table) + { + return hiveColumnHandles(table).stream() + .collect(toImmutableMap(HiveColumnHandle::getName, identity())); + } + @Override public long getTableModificationTime(ConnectorSession session, ConnectorTableHandle tableHandle) { @@ -687,20 +698,23 @@ public class HiveMetadata } @Override - public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint) + public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { if (!HiveSessionProperties.isStatisticsEnabled(session)) { return TableStatistics.empty(); } - Map columns = getColumnHandles(session, tableHandle) + SchemaTableName tableName = ((HiveTableHandle) tableHandle).getSchemaTableName(); + Table table = metastore.getTable(new HiveIdentity(session), tableName.getSchemaName(), tableName.getTableName()) + .orElseThrow(() -> new TableNotFoundException(tableName)); + Map columns = getColumnHandles(table) .entrySet().stream() .filter(entry -> !((HiveColumnHandle) entry.getValue()).isHidden()) .collect(toImmutableMap(Map.Entry::getKey, Map.Entry::getValue)); Map columnTypes = columns.entrySet().stream() .collect(toImmutableMap(Map.Entry::getKey, entry -> getColumnMetadata(session, tableHandle, entry.getValue()).getType())); - HivePartitionResult partitionResult = partitionManager.getPartitions(metastore, new HiveIdentity(session), tableHandle, constraint); + HivePartitionResult partitionResult = partitionManager.getPartitions(metastore, new HiveIdentity(session), tableHandle, constraint, table); List partitions = partitionManager.getPartitionsAsList(partitionResult); - return hiveStatisticsProvider.getTableStatistics(session, ((HiveTableHandle) tableHandle).getSchemaTableName(), columns, columnTypes, partitions); + return hiveStatisticsProvider.getTableStatistics(session, ((HiveTableHandle) tableHandle).getSchemaTableName(), columns, columnTypes, partitions, includeColumnStatistics, table); } private List listTables(ConnectorSession session, SchemaTablePrefix prefix) @@ -2037,8 +2051,11 @@ public class HiveMetadata if (constraint == null) { return Optional.of(handle); } + SchemaTableName tableName = hiveTableHandle.getSchemaTableName(); + Table table = metastore.getTable(new HiveIdentity(session), tableName.getSchemaName(), tableName.getTableName()) + .orElseThrow(() -> new TableNotFoundException(tableName)); HiveIdentity identity = new HiveIdentity(session); - HivePartitionResult partitionResult = partitionManager.getPartitions(metastore, identity, handle, constraint); + HivePartitionResult partitionResult = partitionManager.getPartitions(metastore, identity, handle, constraint, table); HiveTableHandle newHandle = partitionManager.applyPartitionResult(hiveTableHandle, partitionResult); return Optional.of(newHandle); } @@ -2058,7 +2075,7 @@ public class HiveMetadata metastore.truncateUnpartitionedTable(session, handle.getSchemaName(), handle.getTableName()); } else { - for (HivePartition hivePartition : partitionManager.getOrLoadPartitions(metastore, identity, handle)) { + for (HivePartition hivePartition : partitionManager.getOrLoadPartitions(session, metastore, identity, handle)) { metastore.dropPartition(session, handle.getSchemaName(), handle.getTableName(), toPartitionValues(hivePartition.getPartitionId())); } } @@ -2085,7 +2102,7 @@ public class HiveMetadata HiveTableHandle hiveTable = (HiveTableHandle) table; List partitionColumns = ImmutableList.copyOf(hiveTable.getPartitionColumns()); - List partitions = partitionManager.getOrLoadPartitions(metastore, identity, hiveTable); + List partitions = partitionManager.getOrLoadPartitions(session, metastore, identity, hiveTable); TupleDomain predicate = createPredicate(partitionColumns, partitions); @@ -2148,7 +2165,11 @@ public class HiveMetadata HiveTableHandle handle = (HiveTableHandle) tableHandle; checkArgument(!handle.getAnalyzePartitionValues().isPresent() || constraint.getSummary().isAll(), "Analyze should not have a constraint"); - HivePartitionResult partitionResult = partitionManager.getPartitions(metastore, identity, handle, constraint); + SchemaTableName tableName = handle.getSchemaTableName(); + Table table = metastore.getTable(new HiveIdentity(session), tableName.getSchemaName(), tableName.getTableName()) + .orElseThrow(() -> new TableNotFoundException(tableName)); + + HivePartitionResult partitionResult = partitionManager.getPartitions(metastore, identity, handle, constraint, table); HiveTableHandle newHandle = partitionManager.applyPartitionResult(handle, partitionResult); @@ -2204,7 +2225,7 @@ public class HiveMetadata } // Get column handle - Map columnHandles = getColumnHandles(session, handle); + Map columnHandles = getColumnHandles(table); // map predicate columns to hive column handles Map predicateColumns = predicateColumnNames.stream() @@ -2235,8 +2256,6 @@ public class HiveMetadata } if (!pushPartitionsOnly && isSuitableToPush) { - Table table = metastore.getTable(identity, handle.getSchemaName(), handle.getTableName()) - .orElseThrow(() -> new TableNotFoundException(handle.getSchemaTableName())); return Optional.of(new ConstraintApplicationResult<>(newHandle, TupleDomain.all())); } diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveMetadataFactory.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveMetadataFactory.java index d28c5a6c2..f361394f4 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveMetadataFactory.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveMetadataFactory.java @@ -22,12 +22,16 @@ import io.prestosql.plugin.hive.metastore.HiveMetastore; import io.prestosql.plugin.hive.metastore.SemiTransactionalHiveMetastore; import io.prestosql.plugin.hive.security.AccessControlMetadataFactory; import io.prestosql.plugin.hive.statistics.MetastoreHiveStatisticsProvider; +import io.prestosql.plugin.hive.statistics.TableColumnStatistics; import io.prestosql.spi.type.TypeManager; import org.joda.time.DateTimeZone; import javax.inject.Inject; +import java.util.List; +import java.util.Map; import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.ScheduledExecutorService; import java.util.function.Supplier; @@ -39,6 +43,9 @@ public class HiveMetadataFactory { private static final Logger log = Logger.get(HiveMetadataFactory.class); + protected final Map statsCache = new ConcurrentHashMap(); + protected final Map> samplePartitionCache = new ConcurrentHashMap(); + private final boolean allowCorruptWritesForTesting; private final boolean skipDeletionForAlter; private final boolean skipTargetCleanupOnRollback; @@ -213,7 +220,7 @@ public class HiveMetadataFactory partitionUpdateCodec, typeTranslator, prestoVersion, - new MetastoreHiveStatisticsProvider(metastore), + new MetastoreHiveStatisticsProvider(metastore, statsCache, samplePartitionCache), accessControlMetadataFactory.create(metastore), autoVacuumEnabled, vacuumDeltaNumThreshold, diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/HivePartitionManager.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/HivePartitionManager.java index 82f2bf630..bf4a21140 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/HivePartitionManager.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/HivePartitionManager.java @@ -25,6 +25,7 @@ import io.prestosql.plugin.hive.metastore.SemiTransactionalHiveMetastore; import io.prestosql.plugin.hive.metastore.Table; import io.prestosql.spi.PrestoException; import io.prestosql.spi.connector.ColumnHandle; +import io.prestosql.spi.connector.ConnectorSession; import io.prestosql.spi.connector.ConnectorTableHandle; import io.prestosql.spi.connector.Constraint; import io.prestosql.spi.connector.SchemaTableName; @@ -114,7 +115,7 @@ public class HivePartitionManager this.typeManager = requireNonNull(typeManager, "typeManager is null"); } - public HivePartitionResult getPartitions(SemiTransactionalHiveMetastore metastore, HiveIdentity identity, ConnectorTableHandle tableHandle, Constraint constraint) + public HivePartitionResult getPartitions(SemiTransactionalHiveMetastore metastore, HiveIdentity identity, ConnectorTableHandle tableHandle, Constraint constraint, Table table) { HiveTableHandle hiveTableHandle = (HiveTableHandle) tableHandle; TupleDomain effectivePredicate = constraint.getSummary() @@ -128,9 +129,6 @@ public class HivePartitionManager return new HivePartitionResult(partitionColumns, ImmutableList.of(), none(), none(), none(), hiveBucketHandle, Optional.empty()); } - Table table = metastore.getTable(identity, tableName.getSchemaName(), tableName.getTableName()) - .orElseThrow(() -> new TableNotFoundException(tableName)); - Optional bucketFilter = HiveBucketing.getHiveBucketFilter(table, effectivePredicate); TupleDomain compactEffectivePredicate = toCompactTupleDomain(effectivePredicate, domainCompactionThreshold); @@ -157,7 +155,7 @@ public class HivePartitionManager .collect(toImmutableList()); } else { - List partitionNames = getFilteredPartitionNames(metastore, identity, tableName, partitionColumns, effectivePredicate); + List partitionNames = getFilteredPartitionNames(metastore, identity, tableName, partitionColumns, effectivePredicate, table); partitionsIterable = () -> partitionNames.stream() // Apply extra filters which could not be done by getFilteredPartitionNames .map(partitionName -> parseValuesAndFilterPartition(tableName, partitionName, partitionColumns, partitionTypes, effectivePredicate, predicate)) @@ -233,10 +231,13 @@ public class HivePartitionManager handle.isSuitableToPush()); } - public List getOrLoadPartitions(SemiTransactionalHiveMetastore metastore, HiveIdentity identity, HiveTableHandle table) + public List getOrLoadPartitions(ConnectorSession session, SemiTransactionalHiveMetastore metastore, HiveIdentity identity, HiveTableHandle tableHandle) { - return table.getPartitions().orElseGet(() -> - getPartitionsAsList(getPartitions(metastore, identity, table, new Constraint(table.getEnforcedConstraint())))); + SchemaTableName tableName = tableHandle.getSchemaTableName(); + Table table = metastore.getTable(new HiveIdentity(session), tableName.getSchemaName(), tableName.getTableName()) + .orElseThrow(() -> new TableNotFoundException(tableName)); + return tableHandle.getPartitions().orElseGet(() -> + getPartitionsAsList(getPartitions(metastore, identity, tableHandle, new Constraint(tableHandle.getEnforcedConstraint()), table))); } private static TupleDomain toCompactTupleDomain(TupleDomain effectivePredicate, int threshold) @@ -288,7 +289,7 @@ public class HivePartitionManager return constraint.test(partition.getKeys()); } - private List getFilteredPartitionNames(SemiTransactionalHiveMetastore metastore, HiveIdentity identity, SchemaTableName tableName, List partitionKeys, TupleDomain effectivePredicate) + private List getFilteredPartitionNames(SemiTransactionalHiveMetastore metastore, HiveIdentity identity, SchemaTableName tableName, List partitionKeys, TupleDomain effectivePredicate, Table table) { checkArgument(effectivePredicate.getDomains().isPresent()); @@ -351,7 +352,7 @@ public class HivePartitionManager } // fetch the partition names - return metastore.getPartitionNamesByParts(identity, tableName.getSchemaName(), tableName.getTableName(), filter) + return metastore.getPartitionNamesByParts(identity, tableName.getSchemaName(), tableName.getTableName(), filter, table) .orElseThrow(() -> new TableNotFoundException(tableName)); } diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveSplitManager.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveSplitManager.java index b7d8e36af..f33b41485 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveSplitManager.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveSplitManager.java @@ -214,7 +214,7 @@ public class HiveSplitManager } // get partitions - List partitions = partitionManager.getOrLoadPartitions(metastore, new HiveIdentity(session), hiveTable); + List partitions = partitionManager.getOrLoadPartitions(session, metastore, new HiveIdentity(session), hiveTable); // short circuit if we don't have any partitions if (partitions.isEmpty()) { diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveSplitSource.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveSplitSource.java index 1e012b6c0..0c193db58 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveSplitSource.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveSplitSource.java @@ -81,7 +81,7 @@ import static java.util.Objects.requireNonNull; class HiveSplitSource implements ConnectorSplitSource { - private static final Logger log = Logger.get(HiveSplit.class); + private static final Logger log = Logger.get(HiveSplitSource.class); private final String queryId; private final String databaseName; diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveUtil.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveUtil.java index 9d0890406..70ad31f27 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveUtil.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveUtil.java @@ -225,8 +225,8 @@ public final class HiveUtil // Tell hive the columns we would like to read, this lets hive optimize reading column oriented files setReadColumns(configuration, readHiveColumnIndexes); - InputFormat inputFormat = getInputFormat(configuration, schema, true); JobConf jobConf = ConfigurationUtils.toJobConf(configuration); + InputFormat inputFormat = getInputFormat(configuration, schema, true, jobConf); FileSplit fileSplit = new FileSplit(path, start, length, (String[]) null); // propagate serialization configuration to getRecordReader @@ -298,12 +298,10 @@ public final class HiveUtil return Optional.ofNullable(compressionCodecFactory.getCodec(file)); } - static InputFormat getInputFormat(Configuration configuration, Properties schema, boolean symlinkTarget) + static InputFormat getInputFormat(Configuration configuration, Properties schema, boolean symlinkTarget, JobConf jobConf) { String inputFormatName = getInputFormatName(schema); try { - JobConf jobConf = ConfigurationUtils.toJobConf(configuration); - Class> inputFormatClass = getInputFormatClass(jobConf, inputFormatName); if (symlinkTarget && (inputFormatClass == SymlinkTextInputFormat.class)) { // symlink targets are always TextInputFormat diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/metastore/CachingHiveMetastore.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/metastore/CachingHiveMetastore.java index 8d649ba22..f03a51ac9 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/metastore/CachingHiveMetastore.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/metastore/CachingHiveMetastore.java @@ -128,13 +128,14 @@ public class CachingHiveMetastore public static CachingHiveMetastore memoizeMetastore(HiveMetastore delegate, long maximumSize) { + // If delegate is instance of CachingHiveMetastore, we are bypassing directly to second layer of cache, to get cached values. return new CachingHiveMetastore( delegate, newDirectExecutorService(), OptionalLong.empty(), OptionalLong.empty(), maximumSize, - false); + false || delegate instanceof CachingHiveMetastore); } private CachingHiveMetastore(HiveMetastore delegate, Executor executor, OptionalLong expiresAfterWriteMillis, OptionalLong refreshMills, long maximumSize, boolean skipCache) @@ -142,7 +143,8 @@ public class CachingHiveMetastore this.delegate = requireNonNull(delegate, "delegate is null"); requireNonNull(executor, "executor is null"); - this.skipCache = skipCache; + // if refreshMills is present and is 0 , keeps cache unrefreshed. + this.skipCache = skipCache || (refreshMills.isPresent() && refreshMills.getAsLong() == 0); databaseNamesCache = newCacheBuilder(expiresAfterWriteMillis, refreshMills, maximumSize) .build(asyncReloading(CacheLoader.from(this::loadAllDatabases), executor)); @@ -351,9 +353,8 @@ public class CachingHiveMetastore .collect(toImmutableList()); if (skipCache) { - return loadPartitionColumnStatistics(partitions).entrySet() - .stream() - .collect(toImmutableMap(entry -> entry.getKey().getKey().getPartitionName().get(), Entry::getValue)); + HiveIdentity identity1 = updateIdentity(identity); + return delegate.getPartitionStatistics(identity1, table, partitionNames); } Map, PartitionStatistics> statistics = getAll(partitionStatisticsCache, partitions); diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/metastore/SemiTransactionalHiveMetastore.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/metastore/SemiTransactionalHiveMetastore.java index 90a130307..a3547371c 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/metastore/SemiTransactionalHiveMetastore.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/metastore/SemiTransactionalHiveMetastore.java @@ -260,10 +260,9 @@ public class SemiTransactionalHiveMetastore } } - public synchronized Map getPartitionStatistics(HiveIdentity identity, String databaseName, String tableName, Set partitionNames) + public synchronized Map getPartitionStatistics(HiveIdentity identity, String databaseName, String tableName, Set partitionNames, Optional table) { checkReadable(); - Optional
table = getTable(identity, databaseName, tableName); if (!table.isPresent()) { return ImmutableMap.of(); } @@ -606,21 +605,21 @@ public class SemiTransactionalHiveMetastore public synchronized Optional> getPartitionNames(HiveIdentity identity, String databaseName, String tableName) { - return doGetPartitionNames(identity, databaseName, tableName, Optional.empty()); + Optional
table = getTable(identity, databaseName, tableName); + return doGetPartitionNames(identity, databaseName, tableName, Optional.empty(), table); } - public synchronized Optional> getPartitionNamesByParts(HiveIdentity identity, String databaseName, String tableName, List parts) + public synchronized Optional> getPartitionNamesByParts(HiveIdentity identity, String databaseName, String tableName, List parts, Table table) { - return doGetPartitionNames(identity, databaseName, tableName, Optional.of(parts)); + return doGetPartitionNames(identity, databaseName, tableName, Optional.of(parts), Optional.of(table)); } @GuardedBy("this") - private Optional> doGetPartitionNames(HiveIdentity identity, String databaseName, String tableName, Optional> parts) + private Optional> doGetPartitionNames(HiveIdentity identity, String databaseName, String tableName, Optional> parts, Optional
table) { checkHoldsLock(); checkReadable(); - Optional
table = getTable(identity, databaseName, tableName); if (!table.isPresent()) { return Optional.empty(); } diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/HiveStatisticsProvider.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/HiveStatisticsProvider.java index 99e970677..6198e01ab 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/HiveStatisticsProvider.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/HiveStatisticsProvider.java @@ -15,6 +15,7 @@ package io.prestosql.plugin.hive.statistics; import io.prestosql.plugin.hive.HivePartition; +import io.prestosql.plugin.hive.metastore.Table; import io.prestosql.spi.connector.ColumnHandle; import io.prestosql.spi.connector.ConnectorSession; import io.prestosql.spi.connector.SchemaTableName; @@ -31,8 +32,10 @@ public interface HiveStatisticsProvider */ TableStatistics getTableStatistics( ConnectorSession session, - SchemaTableName table, + SchemaTableName schemaTableName, Map columns, Map columnTypes, - List partitions); + List partitions, + boolean includeColumnStatistics, + Table table); } diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/MetastoreHiveStatisticsProvider.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/MetastoreHiveStatisticsProvider.java index 17e717eaf..47d0b3749 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/MetastoreHiveStatisticsProvider.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/MetastoreHiveStatisticsProvider.java @@ -35,6 +35,7 @@ import io.prestosql.plugin.hive.metastore.DoubleStatistics; import io.prestosql.plugin.hive.metastore.HiveColumnStatistics; import io.prestosql.plugin.hive.metastore.IntegerStatistics; import io.prestosql.plugin.hive.metastore.SemiTransactionalHiveMetastore; +import io.prestosql.plugin.hive.metastore.Table; import io.prestosql.spi.PrestoException; import io.prestosql.spi.connector.ColumnHandle; import io.prestosql.spi.connector.ConnectorSession; @@ -96,11 +97,15 @@ public class MetastoreHiveStatisticsProvider private static final Logger log = Logger.get(MetastoreHiveStatisticsProvider.class); private final PartitionsStatisticsProvider statisticsProvider; + private static Map statsCache; + private static Map> samplePartitionCache; - public MetastoreHiveStatisticsProvider(SemiTransactionalHiveMetastore metastore) + public MetastoreHiveStatisticsProvider(SemiTransactionalHiveMetastore metastore, Map statsCache, Map> samplePartitionCache) { requireNonNull(metastore, "metastore is null"); - this.statisticsProvider = (session, table, hivePartitions) -> getPartitionsStatistics(session, metastore, table, hivePartitions); + this.statsCache = requireNonNull(statsCache, "statsCache is null"); + this.samplePartitionCache = requireNonNull(samplePartitionCache, "samplePartitionCache is null"); + this.statisticsProvider = (session, schemaTableName, hivePartitions, table) -> getPartitionsStatistics(session, metastore, schemaTableName, hivePartitions, table); } @VisibleForTesting @@ -109,7 +114,7 @@ public class MetastoreHiveStatisticsProvider this.statisticsProvider = requireNonNull(statisticsProvider, "statisticsProvider is null"); } - private static Map getPartitionsStatistics(ConnectorSession session, SemiTransactionalHiveMetastore metastore, SchemaTableName table, List hivePartitions) + private static Map getPartitionsStatistics(ConnectorSession session, SemiTransactionalHiveMetastore metastore, SchemaTableName schemaTableName, List hivePartitions, Table table) { if (hivePartitions.isEmpty()) { return ImmutableMap.of(); @@ -117,21 +122,23 @@ public class MetastoreHiveStatisticsProvider boolean unpartitioned = hivePartitions.stream().anyMatch(partition -> partition.getPartitionId().equals(UNPARTITIONED_ID)); if (unpartitioned) { checkArgument(hivePartitions.size() == 1, "expected only one hive partition"); - return ImmutableMap.of(UNPARTITIONED_ID, metastore.getTableStatistics(new HiveIdentity(session), table.getSchemaName(), table.getTableName())); + return ImmutableMap.of(UNPARTITIONED_ID, metastore.getTableStatistics(new HiveIdentity(session), schemaTableName.getSchemaName(), schemaTableName.getTableName())); } Set partitionNames = hivePartitions.stream() .map(HivePartition::getPartitionId) .collect(toImmutableSet()); - return metastore.getPartitionStatistics(new HiveIdentity(session), table.getSchemaName(), table.getTableName(), partitionNames); + return metastore.getPartitionStatistics(new HiveIdentity(session), schemaTableName.getSchemaName(), schemaTableName.getTableName(), partitionNames, Optional.of(table)); } @Override public TableStatistics getTableStatistics( ConnectorSession session, - SchemaTableName table, + SchemaTableName schemaTableName, Map columns, Map columnTypes, - List partitions) + List partitions, + boolean includeColumnStatistics, + Table table) { if (!isStatisticsEnabled(session)) { return TableStatistics.empty(); @@ -140,11 +147,25 @@ public class MetastoreHiveStatisticsProvider return createZeroStatistics(columns, columnTypes); } int sampleSize = getPartitionStatisticsSampleSize(session); - List partitionsSample = getPartitionsSample(partitions, sampleSize); + List partitionsSample = samplePartitionCache.get(schemaTableName.getTableName()); + if (includeColumnStatistics || partitionsSample == null) { + partitionsSample = getPartitionsSample(partitions, sampleSize); + samplePartitionCache.put(schemaTableName.getTableName(), partitionsSample); + } try { - Map statisticsSample = statisticsProvider.getPartitionsStatistics(session, table, partitionsSample); - validatePartitionStatistics(table, statisticsSample); - return getTableStatistics(columns, columnTypes, partitions, statisticsSample); + Map statisticsSample = statisticsProvider.getPartitionsStatistics(session, schemaTableName, partitionsSample, table); + if (!includeColumnStatistics) { + OptionalDouble averageRows = calculateAverageRowsPerPartition(statisticsSample.values()); + TableStatistics.Builder result = TableStatistics.builder(); + result.setRowCount(Estimate.of(averageRows.getAsDouble() * partitions.size())); + result.setFileCount(calulateFileCount(statisticsSample.values())); + result.setOnDiskDataSizeInBytes(calculateTotalOnDiskSizeInBytes(statisticsSample.values())); + return result.build(); + } + else { + validatePartitionStatistics(schemaTableName, statisticsSample); + return getTableStatistics(columns, columnTypes, partitions, statisticsSample); + } } catch (PrestoException e) { if (e.getErrorCode().equals(HiveErrorCode.HIVE_CORRUPTED_COLUMN_STATISTICS.toErrorCode()) && isIgnoreCorruptedStatistics(session)) { @@ -404,14 +425,28 @@ public class MetastoreHiveStatisticsProvider double rowCount = averageRowsPerPartition * queriedPartitionsCount; TableStatistics.Builder result = TableStatistics.builder(); + long fileCount = calulateFileCount(statistics.values()); + long totalOnDiskSize = calculateTotalOnDiskSizeInBytes(statistics.values()); result.setRowCount(Estimate.of(rowCount)); + result.setFileCount(fileCount); + result.setOnDiskDataSizeInBytes(totalOnDiskSize); for (Map.Entry column : columns.entrySet()) { String columnName = column.getKey(); HiveColumnHandle columnHandle = (HiveColumnHandle) column.getValue(); Type columnType = columnTypes.get(columnName); ColumnStatistics columnStatistics; + TableColumnStatistics tableColumnStatistics; if (columnHandle.isPartitionKey()) { - columnStatistics = createPartitionColumnStatistics(columnHandle, columnType, partitions, statistics, averageRowsPerPartition, rowCount); + tableColumnStatistics = statsCache.get(partitions.get(0).getTableName().getTableName() + columnName); + if (tableColumnStatistics == null || invalidateStatsCache(partitions.get(0).getTableName().getTableName() + columnName, Estimate.of(rowCount), fileCount, totalOnDiskSize)) { + columnStatistics = createPartitionColumnStatistics(columnHandle, columnType, partitions, statistics, averageRowsPerPartition, rowCount); + TableStatistics tableStatistics = new TableStatistics(Estimate.of(rowCount), fileCount, totalOnDiskSize, ImmutableMap.of()); + tableColumnStatistics = new TableColumnStatistics(tableStatistics, columnStatistics); + statsCache.put(partitions.get(0).getTableName().getTableName() + columnName, tableColumnStatistics); + } + else { + columnStatistics = tableColumnStatistics.columnStatistics; + } } else { columnStatistics = createDataColumnStatistics(columnName, columnType, rowCount, statistics.values()); @@ -421,6 +456,16 @@ public class MetastoreHiveStatisticsProvider return result.build(); } + private static boolean invalidateStatsCache(String tableNameColumName, Estimate rowCount, long fileCount, long totalOnDisk) + { + if (statsCache.get(tableNameColumName).tableStatistics.getOnDiskDataSizeInBytes() != totalOnDisk + || statsCache.get(tableNameColumName).tableStatistics.getFileCount() != fileCount + || !statsCache.get(tableNameColumName).tableStatistics.getRowCount().equals(rowCount)) { + return true; + } + return false; + } + @VisibleForTesting static OptionalDouble calculateAverageRowsPerPartition(Collection statistics) { @@ -433,6 +478,26 @@ public class MetastoreHiveStatisticsProvider .average(); } + static long calulateFileCount(Collection statistics) + { + return statistics.stream() + .map(PartitionStatistics::getBasicStatistics) + .map(HiveBasicStatistics::getFileCount) + .filter(OptionalLong::isPresent) + .mapToLong(OptionalLong::getAsLong) + .sum(); + } + + static long calculateTotalOnDiskSizeInBytes(Collection statistics) + { + return statistics.stream() + .map(PartitionStatistics::getBasicStatistics) + .map(HiveBasicStatistics::getOnDiskDataSizeInBytes) + .filter(OptionalLong::isPresent) + .mapToLong(OptionalLong::getAsLong) + .sum(); + } + private static ColumnStatistics createPartitionColumnStatistics( HiveColumnHandle column, Type type, @@ -847,6 +912,6 @@ public class MetastoreHiveStatisticsProvider @VisibleForTesting interface PartitionsStatisticsProvider { - Map getPartitionsStatistics(ConnectorSession session, SchemaTableName table, List hivePartitions); + Map getPartitionsStatistics(ConnectorSession session, SchemaTableName schemaTableName, List hivePartitions, Table table); } } diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/TableColumnStatistics.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/TableColumnStatistics.java new file mode 100644 index 000000000..43ec12117 --- /dev/null +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/statistics/TableColumnStatistics.java @@ -0,0 +1,30 @@ +/* + * Copyright (C) 2018-2021. Huawei Technologies Co., Ltd. All rights reserved. + * Licensed 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 io.prestosql.plugin.hive.statistics; + +import io.prestosql.spi.statistics.ColumnStatistics; +import io.prestosql.spi.statistics.TableStatistics; + +public class TableColumnStatistics +{ + TableStatistics tableStatistics; + ColumnStatistics columnStatistics; + + public TableColumnStatistics(TableStatistics tableStatistics, ColumnStatistics columnStatistics) + { + this.tableStatistics = tableStatistics; + this.columnStatistics = columnStatistics; + } +} diff --git a/presto-hive/src/test/java/io/prestosql/plugin/hive/AbstractTestHive.java b/presto-hive/src/test/java/io/prestosql/plugin/hive/AbstractTestHive.java index 737de231f..58dda32b7 100644 --- a/presto-hive/src/test/java/io/prestosql/plugin/hive/AbstractTestHive.java +++ b/presto-hive/src/test/java/io/prestosql/plugin/hive/AbstractTestHive.java @@ -1302,7 +1302,7 @@ public abstract class AbstractTestHive ConnectorMetadata metadata = transaction.getMetadata(); ConnectorSession session = newSession(); ConnectorTableHandle tableHandle = getTableHandle(metadata, tableName); - TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, Constraint.alwaysTrue()); + TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, Constraint.alwaysTrue(), true); assertFalse(tableStatistics.getRowCount().isUnknown(), "row count is unknown"); @@ -3075,8 +3075,8 @@ public abstract class AbstractTestHive ConnectorMetadata metadata = transaction.getMetadata(); ConnectorTableHandle tableHandle = metadata.getTableHandle(session, tableName); - TableStatistics unsampledStatistics = metadata.getTableStatistics(sampleSize(2), tableHandle, Constraint.alwaysTrue()); - TableStatistics sampledStatistics = metadata.getTableStatistics(sampleSize(1), tableHandle, Constraint.alwaysTrue()); + TableStatistics unsampledStatistics = metadata.getTableStatistics(sampleSize(2), tableHandle, Constraint.alwaysTrue(), true); + TableStatistics sampledStatistics = metadata.getTableStatistics(sampleSize(1), tableHandle, Constraint.alwaysTrue(), true); assertEquals(sampledStatistics, unsampledStatistics); } } @@ -3923,9 +3923,10 @@ public abstract class AbstractTestHive private static HiveBasicStatistics getBasicStatisticsForPartition(ConnectorSession session, Transaction transaction, SchemaTableName table, String partitionName) { + HiveIdentity identity = new HiveIdentity(session); return transaction .getMetastore(table.getSchemaName()) - .getPartitionStatistics(new HiveIdentity(session), table.getSchemaName(), table.getTableName(), ImmutableSet.of(partitionName)) + .getPartitionStatistics(identity, table.getSchemaName(), table.getTableName(), ImmutableSet.of(partitionName), transaction.getMetastore(table.getSchemaName()).getTable(identity, table.getSchemaName(), table.getTableName())) .get(partitionName) .getBasicStatistics(); } diff --git a/presto-hive/src/test/java/io/prestosql/plugin/hive/TestBackgroundHiveSplitLoader.java b/presto-hive/src/test/java/io/prestosql/plugin/hive/TestBackgroundHiveSplitLoader.java index 111e6c45a..f4fa4e775 100644 --- a/presto-hive/src/test/java/io/prestosql/plugin/hive/TestBackgroundHiveSplitLoader.java +++ b/presto-hive/src/test/java/io/prestosql/plugin/hive/TestBackgroundHiveSplitLoader.java @@ -308,6 +308,7 @@ public class TestBackgroundHiveSplitLoader public void testPropagateException(boolean error, int threads) { AtomicBoolean iteratorUsedAfterException = new AtomicBoolean(); + AtomicBoolean isFirstTime = new AtomicBoolean(true); BackgroundHiveSplitLoader backgroundHiveSplitLoader = new BackgroundHiveSplitLoader( SIMPLE_TABLE, @@ -325,12 +326,19 @@ public class TestBackgroundHiveSplitLoader @Override public HivePartitionMetadata next() { - iteratorUsedAfterException.compareAndSet(false, threw); - threw = true; - if (error) { - throw new Error("loading error occurred"); + // isFirstTime variable is used to skip throwing exception from next method called in BackgroundHiveSplitLoader constructor + if (!isFirstTime.compareAndSet(true, false)) { + iteratorUsedAfterException.compareAndSet(false, threw); + threw = true; + if (error) { + throw new Error("loading error occurred"); + } + throw new RuntimeException("loading error occurred"); } - throw new RuntimeException("loading error occurred"); + return new HivePartitionMetadata( + new HivePartition(new SchemaTableName("testSchema", "table_name")), + Optional.empty(), + ImmutableMap.of()); } }, TupleDomain.all(), diff --git a/presto-hive/src/test/java/io/prestosql/plugin/hive/statistics/TestMetastoreHiveStatisticsProvider.java b/presto-hive/src/test/java/io/prestosql/plugin/hive/statistics/TestMetastoreHiveStatisticsProvider.java index beff4dfb5..019ffeb8c 100644 --- a/presto-hive/src/test/java/io/prestosql/plugin/hive/statistics/TestMetastoreHiveStatisticsProvider.java +++ b/presto-hive/src/test/java/io/prestosql/plugin/hive/statistics/TestMetastoreHiveStatisticsProvider.java @@ -24,11 +24,15 @@ import io.prestosql.plugin.hive.HiveSessionProperties; import io.prestosql.plugin.hive.OrcFileWriterConfig; import io.prestosql.plugin.hive.ParquetFileWriterConfig; import io.prestosql.plugin.hive.PartitionStatistics; +import io.prestosql.plugin.hive.metastore.Column; import io.prestosql.plugin.hive.metastore.DateStatistics; import io.prestosql.plugin.hive.metastore.DecimalStatistics; import io.prestosql.plugin.hive.metastore.DoubleStatistics; import io.prestosql.plugin.hive.metastore.HiveColumnStatistics; import io.prestosql.plugin.hive.metastore.IntegerStatistics; +import io.prestosql.plugin.hive.metastore.Storage; +import io.prestosql.plugin.hive.metastore.StorageFormat; +import io.prestosql.plugin.hive.metastore.Table; import io.prestosql.spi.PrestoException; import io.prestosql.spi.connector.SchemaTableName; import io.prestosql.spi.statistics.ColumnStatistics; @@ -51,6 +55,7 @@ import static io.prestosql.plugin.hive.HiveColumnHandle.ColumnType.PARTITION_KEY import static io.prestosql.plugin.hive.HiveColumnHandle.ColumnType.REGULAR; import static io.prestosql.plugin.hive.HivePartition.UNPARTITIONED_ID; import static io.prestosql.plugin.hive.HivePartitionManager.parsePartition; +import static io.prestosql.plugin.hive.HiveStorageFormat.ORC; import static io.prestosql.plugin.hive.HiveType.HIVE_LONG; import static io.prestosql.plugin.hive.HiveType.HIVE_STRING; import static io.prestosql.plugin.hive.HiveUtil.parsePartitionValue; @@ -84,6 +89,7 @@ import static org.testng.Assert.assertEquals; public class TestMetastoreHiveStatisticsProvider { + private static final Storage STORAGE_1 = new Storage(StorageFormat.fromHiveStorageFormat(ORC), "", Optional.empty(), false, ImmutableMap.of()); private static final SchemaTableName TABLE = new SchemaTableName("schema", "table"); private static final String PARTITION = "partition"; private static final String COLUMN = "column"; @@ -91,6 +97,7 @@ public class TestMetastoreHiveStatisticsProvider private static final HiveColumnHandle PARTITION_COLUMN_1 = new HiveColumnHandle("p1", HIVE_STRING, VARCHAR.getTypeSignature(), 0, PARTITION_KEY, Optional.empty()); private static final HiveColumnHandle PARTITION_COLUMN_2 = new HiveColumnHandle("p2", HIVE_LONG, BIGINT.getTypeSignature(), 1, PARTITION_KEY, Optional.empty()); + private static final Table table = new Table(TABLE.getSchemaName(), TABLE.getTableName(), "user", "MANAGED_TABLE", STORAGE_1, ImmutableList.of(), ImmutableList.of(new Column("p1", HIVE_STRING, Optional.empty()), new Column("p2", HIVE_LONG, Optional.empty())), ImmutableMap.of(), Optional.of("original"), Optional.of("expanded")); @Test public void testGetPartitionsSample() @@ -604,7 +611,7 @@ public class TestMetastoreHiveStatisticsProvider .setBasicStatistics(new HiveBasicStatistics(OptionalLong.empty(), OptionalLong.of(1000), OptionalLong.empty(), OptionalLong.empty())) .setColumnStatistics(ImmutableMap.of(COLUMN, HiveColumnStatistics.createIntegerColumnStatistics(OptionalLong.of(-100), OptionalLong.of(100), OptionalLong.of(500), OptionalLong.of(300)))) .build(); - MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, table, hivePartitions) -> ImmutableMap.of(partitionName, statistics)); + MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, schemaTableName, hivePartitions, table) -> ImmutableMap.of(partitionName, statistics)); TestingConnectorSession session = new TestingConnectorSession(new HiveSessionProperties(new HiveConfig(), new OrcFileWriterConfig(), new ParquetFileWriterConfig()).getSessionProperties()); HiveColumnHandle columnHandle = new HiveColumnHandle(COLUMN, HIVE_LONG, BIGINT.getTypeSignature(), 2, REGULAR, Optional.empty()); TableStatistics expected = TableStatistics.builder() @@ -643,7 +650,7 @@ public class TestMetastoreHiveStatisticsProvider "p1", VARCHAR, "p2", BIGINT, COLUMN, BIGINT), - ImmutableList.of(partition(partitionName))), + ImmutableList.of(partition(partitionName)), true, table), expected); } @@ -654,7 +661,7 @@ public class TestMetastoreHiveStatisticsProvider .setBasicStatistics(new HiveBasicStatistics(OptionalLong.empty(), OptionalLong.of(1000), OptionalLong.empty(), OptionalLong.empty())) .setColumnStatistics(ImmutableMap.of(COLUMN, HiveColumnStatistics.createIntegerColumnStatistics(OptionalLong.of(-100), OptionalLong.of(100), OptionalLong.of(500), OptionalLong.of(300)))) .build(); - MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, table, hivePartitions) -> ImmutableMap.of(UNPARTITIONED_ID, statistics)); + MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, schemaTableName, hivePartitions, table) -> ImmutableMap.of(UNPARTITIONED_ID, statistics)); TestingConnectorSession session = new TestingConnectorSession(new HiveSessionProperties(new HiveConfig(), new OrcFileWriterConfig(), new ParquetFileWriterConfig()).getSessionProperties()); HiveColumnHandle columnHandle = new HiveColumnHandle(COLUMN, HIVE_LONG, BIGINT.getTypeSignature(), 2, REGULAR, Optional.empty()); TableStatistics expected = TableStatistics.builder() @@ -673,7 +680,7 @@ public class TestMetastoreHiveStatisticsProvider TABLE, ImmutableMap.of(COLUMN, columnHandle), ImmutableMap.of(COLUMN, BIGINT), - ImmutableList.of(new HivePartition(TABLE))), + ImmutableList.of(new HivePartition(TABLE)), true, table), expected); } @@ -681,7 +688,7 @@ public class TestMetastoreHiveStatisticsProvider public void testGetTableStatisticsEmpty() { String partitionName = "p1=string1/p2=1234"; - MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, table, hivePartitions) -> ImmutableMap.of(partitionName, PartitionStatistics.empty())); + MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, schemaTableName, hivePartitions, table) -> ImmutableMap.of(partitionName, PartitionStatistics.empty())); TestingConnectorSession session = new TestingConnectorSession(new HiveSessionProperties(new HiveConfig(), new OrcFileWriterConfig(), new ParquetFileWriterConfig()).getSessionProperties()); assertEquals( statisticsProvider.getTableStatistics( @@ -689,15 +696,15 @@ public class TestMetastoreHiveStatisticsProvider TABLE, ImmutableMap.of(), ImmutableMap.of(), - ImmutableList.of(partition(partitionName))), + ImmutableList.of(partition(partitionName)), true, table), TableStatistics.empty()); } @Test public void testGetTableStatisticsSampling() { - MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, table, hivePartitions) -> { - assertEquals(table, TABLE); + MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, schemaTableName, hivePartitions, table) -> { + assertEquals(schemaTableName, TABLE); assertEquals(hivePartitions.size(), 1); return ImmutableMap.of(); }); @@ -711,7 +718,7 @@ public class TestMetastoreHiveStatisticsProvider TABLE, ImmutableMap.of(), ImmutableMap.of(), - ImmutableList.of(partition("p1=string1/p2=1234"), partition("p1=string1/p2=1235"))); + ImmutableList.of(partition("p1=string1/p2=1234"), partition("p1=string1/p2=1235")), true, table); } @Test @@ -721,7 +728,7 @@ public class TestMetastoreHiveStatisticsProvider .setBasicStatistics(new HiveBasicStatistics(-1, 0, 0, 0)) .build(); String partitionName = "p1=string1/p2=1234"; - MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, table, hivePartitions) -> ImmutableMap.of(partitionName, corruptedStatistics)); + MetastoreHiveStatisticsProvider statisticsProvider = new MetastoreHiveStatisticsProvider((session, schemaTableName, hivePartitions, table) -> ImmutableMap.of(partitionName, corruptedStatistics)); TestingConnectorSession session = new TestingConnectorSession(new HiveSessionProperties( new HiveConfig().setIgnoreCorruptedStatistics(false), new OrcFileWriterConfig(), @@ -732,7 +739,7 @@ public class TestMetastoreHiveStatisticsProvider TABLE, ImmutableMap.of(), ImmutableMap.of(), - ImmutableList.of(partition(partitionName)))) + ImmutableList.of(partition(partitionName)), true, table)) .isInstanceOf(PrestoException.class) .hasFieldOrPropertyWithValue("errorCode", HiveErrorCode.HIVE_CORRUPTED_COLUMN_STATISTICS.toErrorCode()); TestingConnectorSession ignoreSession = new TestingConnectorSession(new HiveSessionProperties( @@ -746,7 +753,7 @@ public class TestMetastoreHiveStatisticsProvider TABLE, ImmutableMap.of(), ImmutableMap.of(), - ImmutableList.of(partition(partitionName))), + ImmutableList.of(partition(partitionName)), true, table), TableStatistics.empty()); } diff --git a/presto-main/src/main/java/io/prestosql/cost/TableScanStatsRule.java b/presto-main/src/main/java/io/prestosql/cost/TableScanStatsRule.java index 3e7db6396..ccbb77ee6 100644 --- a/presto-main/src/main/java/io/prestosql/cost/TableScanStatsRule.java +++ b/presto-main/src/main/java/io/prestosql/cost/TableScanStatsRule.java @@ -73,7 +73,7 @@ public class TableScanStatsRule TupleDomain predicate = metadata.getTableProperties(session, node.getTable()).getPredicate(); Constraint constraint = new Constraint(predicate); - TableStatistics tableStatistics = metadata.getTableStatistics(session, node.getTable(), constraint); + TableStatistics tableStatistics = metadata.getTableStatistics(session, node.getTable(), constraint, true); verify(tableStatistics != null, "tableStatistics is null for %s", node); Map outputSymbolStats = new HashMap<>(); diff --git a/presto-main/src/main/java/io/prestosql/dispatcher/LocalDispatchQueryFactory.java b/presto-main/src/main/java/io/prestosql/dispatcher/LocalDispatchQueryFactory.java index 39ee3ab5d..2a69d52f4 100644 --- a/presto-main/src/main/java/io/prestosql/dispatcher/LocalDispatchQueryFactory.java +++ b/presto-main/src/main/java/io/prestosql/dispatcher/LocalDispatchQueryFactory.java @@ -113,12 +113,15 @@ public class LocalDispatchQueryFactory queryMonitor.queryCreatedEvent(stateMachine.getBasicQueryInfo(Optional.empty())); ListenableFuture queryExecutionFuture = executor.submit(() -> { + stateMachine.beginSyntaxAnalysis(); QueryExecutionFactory queryExecutionFactory = executionFactories.get(preparedQuery.getStatement().getClass()); if (queryExecutionFactory == null) { throw new PrestoException(NOT_SUPPORTED, "Unsupported statement type: " + preparedQuery.getStatement().getClass().getSimpleName()); } - return queryExecutionFactory.createQueryExecution(preparedQuery, stateMachine, slug, warningCollector); + QueryExecution queryExecution = queryExecutionFactory.createQueryExecution(preparedQuery, stateMachine, slug, warningCollector); + stateMachine.endSyntaxAnalysis(); + return queryExecution; }); return new LocalDispatchQuery( diff --git a/presto-main/src/main/java/io/prestosql/event/QueryMonitor.java b/presto-main/src/main/java/io/prestosql/event/QueryMonitor.java index a6d2551b2..ccf5a8804 100644 --- a/presto-main/src/main/java/io/prestosql/event/QueryMonitor.java +++ b/presto-main/src/main/java/io/prestosql/event/QueryMonitor.java @@ -402,6 +402,10 @@ public class QueryMonitor // planning duration -- start to end of planning long planning = queryStats.getTotalPlanningTime().toMillis(); + long logicalPlanning = queryStats.getTotalLogicalPlanningTime().toMillis(); + long distributedPlanning = queryStats.getDistributedPlanningTime().toMillis(); + long physicalPlanning = queryStats.getAnalysisTime().toMillis() - logicalPlanning; + long syntaxAnalysisTime = queryStats.getTotalSyntaxAnalysisTime().toMillis(); // Time spent waiting for required no. of worker nodes to be present long waiting = queryStats.getResourceWaitingTime().toMillis(); @@ -446,7 +450,11 @@ public class QueryMonitor queryInfo.getQueryId(), queryInfo.getSession().getTransactionId().map(TransactionId::toString).orElse(""), elapsed, + syntaxAnalysisTime, planning, + logicalPlanning, + physicalPlanning, + distributedPlanning, waiting, scheduling, running, @@ -475,11 +483,15 @@ public class QueryMonitor queryInfo.getQueryId(), queryInfo.getSession().getTransactionId().map(TransactionId::toString).orElse(""), elapsed, + 0, elapsed, 0, 0, 0, 0, + 0, + 0, + 0, queryStartTime, queryEndTime); } @@ -488,7 +500,11 @@ public class QueryMonitor QueryId queryId, String transactionId, long elapsedMillis, + long syntaxAnalysisTime, long planningMillis, + long logicalPlanningMillis, + long physicalPlanningMillis, + long distributedPlanningMillis, long waitingMillis, long schedulingMillis, long runningMillis, @@ -496,13 +512,17 @@ public class QueryMonitor DateTime queryStartTime, DateTime queryEndTime) { - log.info("TIMELINE: Query %s :: Transaction:[%s] :: elapsed %sms :: planning %sms :: waiting %sms :: scheduling %sms :: running %sms :: finishing %sms :: begin %s :: end %s", + log.info("TIMELINE: Query %s :: Transaction:[%s] :: elapsed %sms :: syntaxAnalysisTime %sms :: planning %sms :: logicalPlanningMillis %sms :: physicalPlanningMillis %sms :: distributionPlanTime %sms :: waiting %sms :: scheduling %sms :: running %sms :: finishing %sms :: begin %s :: end %s", queryId, transactionId, elapsedMillis, + syntaxAnalysisTime, planningMillis, - waitingMillis, - schedulingMillis, + logicalPlanningMillis, + physicalPlanningMillis, + distributedPlanningMillis, + (waitingMillis - syntaxAnalysisTime) < 0 ? 0 : waitingMillis - syntaxAnalysisTime, + schedulingMillis - waitingMillis, runningMillis, finishingMillis, queryStartTime, diff --git a/presto-main/src/main/java/io/prestosql/execution/QueryStateMachine.java b/presto-main/src/main/java/io/prestosql/execution/QueryStateMachine.java index d986c4d31..ae032d0d7 100644 --- a/presto-main/src/main/java/io/prestosql/execution/QueryStateMachine.java +++ b/presto-main/src/main/java/io/prestosql/execution/QueryStateMachine.java @@ -106,6 +106,7 @@ public class QueryStateMachine private final URI self; private final ResourceGroupId resourceGroup; private final ResourceGroupManager resourceGroupManager; + private boolean throttlingEnabled; private final TransactionManager transactionManager; private final Metadata metadata; private final QueryOutputManager outputManager; @@ -178,6 +179,8 @@ public class QueryStateMachine this.self = requireNonNull(self, "self is null"); this.resourceGroup = requireNonNull(resourceGroup, "resourceGroup is null"); this.resourceGroupManager = resourceGroupManager; + this.throttlingEnabled = resourceGroupManager.isGroupRegistered(resourceGroup) + && resourceGroupManager.getSoftReservedMemory(resourceGroup) != Long.MAX_VALUE; this.transactionManager = requireNonNull(transactionManager, "transactionManager is null"); this.queryStateTimer = new QueryStateTimer(ticker); this.metadata = requireNonNull(metadata, "metadata is null"); @@ -275,6 +278,11 @@ public class QueryStateMachine return resourceGroupManager; } + public boolean isThrottlingEnabled() + { + return throttlingEnabled; + } + public Session getSession() { return session; @@ -562,6 +570,8 @@ public class QueryStateMachine queryStateTimer.getAnalysisTime(), queryStateTimer.getDistributedPlanningTime(), queryStateTimer.getPlanningTime(), + queryStateTimer.getLogicalPlanningTime(), + queryStateTimer.getSyntaxAnalysisTime(), queryStateTimer.getFinishingTime(), totalTasks, @@ -968,6 +978,16 @@ public class QueryStateMachine queryStateTimer.recordHeartbeat(); } + public void beginSyntaxAnalysis() + { + queryStateTimer.beginSyntaxAnalysis(); + } + + public void endSyntaxAnalysis() + { + queryStateTimer.endSyntaxAnalysis(); + } + public void beginAnalysis() { queryStateTimer.beginAnalyzing(); @@ -978,6 +998,16 @@ public class QueryStateMachine queryStateTimer.endAnalysis(); } + public void beginLogicalPlan() + { + queryStateTimer.beginLogicalPlan(); + } + + public void endLogicalPlan() + { + queryStateTimer.endLogicalPlan(); + } + public void beginDistributedPlanning() { queryStateTimer.beginDistributedPlanning(); @@ -1127,6 +1157,8 @@ public class QueryStateMachine queryStats.getAnalysisTime(), queryStats.getDistributedPlanningTime(), queryStats.getTotalPlanningTime(), + queryStats.getTotalLogicalPlanningTime(), + queryStats.getTotalSyntaxAnalysisTime(), queryStats.getFinishingTime(), queryStats.getTotalTasks(), queryStats.getRunningTasks(), diff --git a/presto-main/src/main/java/io/prestosql/execution/QueryStateTimer.java b/presto-main/src/main/java/io/prestosql/execution/QueryStateTimer.java index 7f1a2f1ab..0fa5bfab9 100644 --- a/presto-main/src/main/java/io/prestosql/execution/QueryStateTimer.java +++ b/presto-main/src/main/java/io/prestosql/execution/QueryStateTimer.java @@ -48,6 +48,11 @@ class QueryStateTimer private final AtomicReference beginAnalysisNanos = new AtomicReference<>(); private final AtomicReference analysisTime = new AtomicReference<>(); + private final AtomicReference beginSyntaxAnalysisNanos = new AtomicReference<>(); + private final AtomicReference syntaxAnalysisTime = new AtomicReference<>(); + + private final AtomicReference beginLogicalPlanNanos = new AtomicReference<>(); + private final AtomicReference logicalPlanTime = new AtomicReference<>(); private final AtomicReference beginDistributedPlanningNanos = new AtomicReference<>(); private final AtomicReference distributedPlanningTime = new AtomicReference<>(); @@ -149,6 +154,16 @@ class QueryStateTimer // Additional timings // + public void beginSyntaxAnalysis() + { + beginSyntaxAnalysisNanos.compareAndSet(null, tickerNanos()); + } + + public void endSyntaxAnalysis() + { + syntaxAnalysisTime.compareAndSet(null, nanosSince(beginSyntaxAnalysisNanos, tickerNanos())); + } + public void beginAnalyzing() { beginAnalysisNanos.compareAndSet(null, tickerNanos()); @@ -159,6 +174,16 @@ class QueryStateTimer analysisTime.compareAndSet(null, nanosSince(beginAnalysisNanos, tickerNanos())); } + public void beginLogicalPlan() + { + beginLogicalPlanNanos.compareAndSet(null, tickerNanos()); + } + + public void endLogicalPlan() + { + logicalPlanTime.compareAndSet(null, nanosSince(beginLogicalPlanNanos, tickerNanos())); + } + public void beginDistributedPlanning() { beginDistributedPlanningNanos.compareAndSet(null, tickerNanos()); @@ -222,6 +247,11 @@ class QueryStateTimer return getDuration(planningTime, beginPlanningNanos); } + public Duration getLogicalPlanningTime() + { + return getDuration(logicalPlanTime, beginLogicalPlanNanos); + } + public Duration getFinishingTime() { return getDuration(finishingTime, beginFinishingNanos); @@ -237,6 +267,11 @@ class QueryStateTimer return toDateTime(endNanos); } + public Duration getSyntaxAnalysisTime() + { + return getDuration(syntaxAnalysisTime, beginSyntaxAnalysisNanos); + } + public Duration getAnalysisTime() { return getDuration(analysisTime, beginAnalysisNanos); diff --git a/presto-main/src/main/java/io/prestosql/execution/QueryStats.java b/presto-main/src/main/java/io/prestosql/execution/QueryStats.java index 54e0b43ab..bc34fac54 100644 --- a/presto-main/src/main/java/io/prestosql/execution/QueryStats.java +++ b/presto-main/src/main/java/io/prestosql/execution/QueryStats.java @@ -52,6 +52,8 @@ public class QueryStats private final Duration analysisTime; private final Duration distributedPlanningTime; private final Duration totalPlanningTime; + private final Duration totalLogicalPlanningTime; + private final Duration totalSyntaxAnalysisTime; private final Duration finishingTime; private final int totalTasks; @@ -118,6 +120,8 @@ public class QueryStats @JsonProperty("analysisTime") Duration analysisTime, @JsonProperty("distributedPlanningTime") Duration distributedPlanningTime, @JsonProperty("totalPlanningTime") Duration totalPlanningTime, + @JsonProperty("totalLogicalPlanningTime") Duration totalLogicalPlanningTime, + @JsonProperty("totalSyntaxAnalysisTime") Duration totalSyntaxAnalysisTime, @JsonProperty("finishingTime") Duration finishingTime, @JsonProperty("totalTasks") int totalTasks, @@ -182,6 +186,8 @@ public class QueryStats this.analysisTime = requireNonNull(analysisTime, "analysisTime is null"); this.distributedPlanningTime = requireNonNull(distributedPlanningTime, "distributedPlanningTime is null"); this.totalPlanningTime = requireNonNull(totalPlanningTime, "totalPlanningTime is null"); + this.totalLogicalPlanningTime = requireNonNull(totalLogicalPlanningTime, "totalLogicalPlanningTime is null"); + this.totalSyntaxAnalysisTime = requireNonNull(totalSyntaxAnalysisTime, "totalSyntaxAnalysisTime is null"); this.finishingTime = requireNonNull(finishingTime, "finishingTime is null"); checkArgument(totalTasks >= 0, "totalTasks is negative"); @@ -319,6 +325,18 @@ public class QueryStats return totalPlanningTime; } + @JsonProperty + public Duration getTotalLogicalPlanningTime() + { + return totalLogicalPlanningTime; + } + + @JsonProperty + public Duration getTotalSyntaxAnalysisTime() + { + return totalSyntaxAnalysisTime; + } + @JsonProperty public Duration getFinishingTime() { diff --git a/presto-main/src/main/java/io/prestosql/execution/SqlQueryExecution.java b/presto-main/src/main/java/io/prestosql/execution/SqlQueryExecution.java index 4e7245bc0..d48efd9cc 100644 --- a/presto-main/src/main/java/io/prestosql/execution/SqlQueryExecution.java +++ b/presto-main/src/main/java/io/prestosql/execution/SqlQueryExecution.java @@ -660,6 +660,7 @@ public class SqlQueryExecution { // time analysis phase stateMachine.beginAnalysis(); + stateMachine.beginLogicalPlan(); // plan query PlanNodeIdAllocator idAllocator = new PlanNodeIdAllocator(); @@ -672,6 +673,7 @@ public class SqlQueryExecution // extract output stateMachine.setOutput(analysis.getTarget()); + stateMachine.endLogicalPlan(); // fragment the plan SubPlan fragmentedPlan = planFragmenter.createSubPlans(stateMachine.getSession(), plan, false, stateMachine.getWarningCollector()); diff --git a/presto-main/src/main/java/io/prestosql/execution/resourcegroups/InternalResourceGroupManager.java b/presto-main/src/main/java/io/prestosql/execution/resourcegroups/InternalResourceGroupManager.java index 3e7768d48..6527fb091 100644 --- a/presto-main/src/main/java/io/prestosql/execution/resourcegroups/InternalResourceGroupManager.java +++ b/presto-main/src/main/java/io/prestosql/execution/resourcegroups/InternalResourceGroupManager.java @@ -394,4 +394,16 @@ public final class InternalResourceGroupManager { return groups.get(resourceGroupId).getCachedMemoryUsageBytes(); } + + @Override + public long getSoftReservedMemory(ResourceGroupId resourceGroupId) + { + return groups.get(resourceGroupId).getSoftReservedMemory().toBytes(); + } + + @Override + public boolean isGroupRegistered(ResourceGroupId resourceGroupId) + { + return groups.containsKey(resourceGroupId); + } } diff --git a/presto-main/src/main/java/io/prestosql/execution/resourcegroups/ResourceGroupManager.java b/presto-main/src/main/java/io/prestosql/execution/resourcegroups/ResourceGroupManager.java index c024c93b7..9de52028a 100644 --- a/presto-main/src/main/java/io/prestosql/execution/resourcegroups/ResourceGroupManager.java +++ b/presto-main/src/main/java/io/prestosql/execution/resourcegroups/ResourceGroupManager.java @@ -49,4 +49,14 @@ public interface ResourceGroupManager { return 0; } + + default long getSoftReservedMemory(ResourceGroupId resourceGroupId) + { + return Long.MAX_VALUE; + } + + default boolean isGroupRegistered(ResourceGroupId resourceGroupId) + { + return false; + } } diff --git a/presto-main/src/main/java/io/prestosql/execution/scheduler/SqlQueryScheduler.java b/presto-main/src/main/java/io/prestosql/execution/scheduler/SqlQueryScheduler.java index 3b3c34df3..65951d9bf 100644 --- a/presto-main/src/main/java/io/prestosql/execution/scheduler/SqlQueryScheduler.java +++ b/presto-main/src/main/java/io/prestosql/execution/scheduler/SqlQueryScheduler.java @@ -720,7 +720,7 @@ public class SqlQueryScheduler private boolean canScheduleMoreSplits() { long cachedMemoryUsage = queryStateMachine.getResourceGroupManager().getCachedMemoryUsage(queryStateMachine.getResourceGroup()); - long softReservedMemory = queryStateMachine.getResourceGroupManager().getResourceGroupInfo(queryStateMachine.getResourceGroup()).getSoftReservedMemory().toBytes(); + long softReservedMemory = queryStateMachine.getResourceGroupManager().getSoftReservedMemory(queryStateMachine.getResourceGroup()); if (cachedMemoryUsage < softReservedMemory) { return true; } @@ -747,7 +747,7 @@ public class SqlQueryScheduler // configured limit. If yes throttle further split scheduling. // Throttle Logic: Wait for x seconds (Wait time will increase till max as per THROTTLE_SLEEP_TIMER) // and then let it schedule 10% of splits. - if (!canScheduleMoreSplits()) { + if (queryStateMachine.isThrottlingEnabled() && !canScheduleMoreSplits()) { try { SECONDS.sleep(THROTTLE_SLEEP_TIMER[currentTimerLevel]); } diff --git a/presto-main/src/main/java/io/prestosql/metadata/Metadata.java b/presto-main/src/main/java/io/prestosql/metadata/Metadata.java index 830862e1d..1cbc22a01 100755 --- a/presto-main/src/main/java/io/prestosql/metadata/Metadata.java +++ b/presto-main/src/main/java/io/prestosql/metadata/Metadata.java @@ -104,9 +104,9 @@ public interface Metadata TableMetadata getTableMetadata(Session session, TableHandle tableHandle); /** - * Return statistics for specified table for given filtering contraint. + * Return statistics for specified table for given filtering contraint with a check either to include ColumnStatistics or not */ - TableStatistics getTableStatistics(Session session, TableHandle tableHandle, Constraint constraint); + TableStatistics getTableStatistics(Session session, TableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics); /** * Get the names that match the specified table prefix (never null). diff --git a/presto-main/src/main/java/io/prestosql/metadata/MetadataManager.java b/presto-main/src/main/java/io/prestosql/metadata/MetadataManager.java index 38236e333..d50620918 100755 --- a/presto-main/src/main/java/io/prestosql/metadata/MetadataManager.java +++ b/presto-main/src/main/java/io/prestosql/metadata/MetadataManager.java @@ -371,7 +371,6 @@ public final class MetadataManager .get() .getTableProperties()); } - return new TableProperties(catalogName, handle.getTransaction(), metadata.getTableProperties(connectorSession, handle.getConnectorHandle())); } @@ -448,11 +447,11 @@ public final class MetadataManager } @Override - public TableStatistics getTableStatistics(Session session, TableHandle tableHandle, Constraint constraint) + public TableStatistics getTableStatistics(Session session, TableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { CatalogName catalogName = tableHandle.getCatalogName(); ConnectorMetadata metadata = getMetadata(session, catalogName); - return metadata.getTableStatistics(session.toConnectorSession(catalogName), tableHandle.getConnectorHandle(), constraint); + return metadata.getTableStatistics(session.toConnectorSession(catalogName), tableHandle.getConnectorHandle(), constraint, includeColumnStatistics); } @Override diff --git a/presto-main/src/main/java/io/prestosql/query/CachedSqlQueryExecution.java b/presto-main/src/main/java/io/prestosql/query/CachedSqlQueryExecution.java index 687476ce8..cac3e0dfd 100644 --- a/presto-main/src/main/java/io/prestosql/query/CachedSqlQueryExecution.java +++ b/presto-main/src/main/java/io/prestosql/query/CachedSqlQueryExecution.java @@ -213,9 +213,14 @@ public class CachedSqlQueryExecution cachedPlan.getStatement().equals(statement) && session.getTransactionId().isPresent() && cachedPlan.getIdentity().getUser().equals(session.getIdentity().getUser())) { // TODO: traverse the statement and accept partial match root = plan.getRoot(); try { - if (!cachedPlan.getTableStatistics().equals(tableStatistics)) { - // TableStatistics have changed, therefore the cached plan may no longer be applicable - throw new NoSuchElementException(); + if (!isEqualBasicStatistics(cachedPlan.getTableStatistics(), tableStatistics, tableNames)) { + for (TableHandle tableHandle : analysis.getTables()) { + tableStatistics.replace(tableHandle.getFullyQualifiedName(), metadata.getTableStatistics(session, tableHandle, Constraint.alwaysTrue(), true)); + } + if (!cachedPlan.getTableStatistics().equals(tableStatistics)) { + // TableStatistics have changed, therefore the cached plan may no longer be applicable + throw new NoSuchElementException(); + } } // TableScanNode may contain the old transaction id. // The following logic rewrites the logical plan by replacing the TableScanNode with a new TableScanNode which @@ -233,6 +238,9 @@ public class CachedSqlQueryExecution } else { // Build a new plan + for (TableHandle tableHandle : analysis.getTables()) { + tableStatistics.replace(tableHandle.getFullyQualifiedName(), metadata.getTableStatistics(session, tableHandle, Constraint.alwaysTrue(), true)); + } plan = createAndCachePlan(key, logicalPlanner, statement, tableNames, tableStatistics, optimizers, analysis, columnTypes, systemSessionProperties); root = plan.getRoot(); } @@ -278,7 +286,8 @@ public class CachedSqlQueryExecution try { if (metadata.isExecutionPlanCacheSupported(session, tableHandle)) { tables.add(tableHandle.getFullyQualifiedName()); - tableStatistics.put(tableHandle.getFullyQualifiedName(), metadata.getTableStatistics(session, tableHandle, Constraint.alwaysTrue())); // TODO: Find a way to get constraints instead of reading all table statistics + // includeColumnStatistics is passed as false, so that calculation of columnStatistics are skipped + tableStatistics.put(tableHandle.getFullyQualifiedName(), metadata.getTableStatistics(session, tableHandle, Constraint.alwaysTrue(), false)); // TODO: Find a way to get constraints instead of reading all table statistics Map columnHandles = metadata.getColumnHandles(session, tableHandle); for (ColumnHandle columnHandle : columnHandles.values()) { @@ -298,6 +307,22 @@ public class CachedSqlQueryExecution return true; } + private boolean isEqualBasicStatistics(Map cacheTableStatistics, Map tableStatistics, List tableNames) + { + for (String tableName : tableNames) { + TableStatistics cacheTableStatisticsTemp = cacheTableStatistics.get(tableName); + TableStatistics tableStatisticsTemp = tableStatistics.get(tableName); + if (cacheTableStatisticsTemp == null || + tableStatisticsTemp == null || + cacheTableStatisticsTemp.getFileCount() != tableStatisticsTemp.getFileCount() || + !cacheTableStatisticsTemp.getRowCount().equals(tableStatisticsTemp.getRowCount()) || + cacheTableStatisticsTemp.getOnDiskDataSizeInBytes() != tableStatisticsTemp.getOnDiskDataSizeInBytes()) { + return false; + } + } + return true; + } + private boolean isCacheable(Statement statement) { // Skip cache when creating tables, hack for outdated metadata diff --git a/presto-main/src/main/java/io/prestosql/queryeditorui/execution/Execution.java b/presto-main/src/main/java/io/prestosql/queryeditorui/execution/Execution.java index f04f8a708..117e782f5 100644 --- a/presto-main/src/main/java/io/prestosql/queryeditorui/execution/Execution.java +++ b/presto-main/src/main/java/io/prestosql/queryeditorui/execution/Execution.java @@ -332,6 +332,8 @@ public class Execution zeroDuration, zeroDuration, zeroDuration, + zeroDuration, + zeroDuration, 0, 0, 0, diff --git a/presto-main/src/main/java/io/prestosql/server/QueryResource.java b/presto-main/src/main/java/io/prestosql/server/QueryResource.java index 4bd90a4b6..8baf855bb 100644 --- a/presto-main/src/main/java/io/prestosql/server/QueryResource.java +++ b/presto-main/src/main/java/io/prestosql/server/QueryResource.java @@ -334,6 +334,8 @@ public class QueryResource ZERO_MILLIS, ZERO_MILLIS, ZERO_MILLIS, + ZERO_MILLIS, + ZERO_MILLIS, 0, 0, 0, diff --git a/presto-main/src/main/java/io/prestosql/sql/planner/iterative/rule/RemoveUnsupportedDynamicFilters.java b/presto-main/src/main/java/io/prestosql/sql/planner/iterative/rule/RemoveUnsupportedDynamicFilters.java index 05762332f..fbec451e5 100644 --- a/presto-main/src/main/java/io/prestosql/sql/planner/iterative/rule/RemoveUnsupportedDynamicFilters.java +++ b/presto-main/src/main/java/io/prestosql/sql/planner/iterative/rule/RemoveUnsupportedDynamicFilters.java @@ -308,7 +308,7 @@ public class RemoveUnsupportedDynamicFilters return true; } - Estimate totalRowCount = metadata.getTableStatistics(session, ((TableScanNode) buildSideTableScanNode.get()).getTable(), Constraint.alwaysTrue()).getRowCount(); + Estimate totalRowCount = metadata.getTableStatistics(session, ((TableScanNode) buildSideTableScanNode.get()).getTable(), Constraint.alwaysTrue(), true).getRowCount(); PlanNodeStatsEstimate filteredStats = statsProvider.getStats(node); if (!filteredStats.isOutputRowCountUnknown() && !totalRowCount.isUnknown()) { @@ -328,7 +328,7 @@ public class RemoveUnsupportedDynamicFilters private boolean highSelectivity(FilterNode node) { - Estimate totalRowCount = metadata.getTableStatistics(session, ((TableScanNode) node.getSource()).getTable(), Constraint.alwaysTrue()).getRowCount(); + Estimate totalRowCount = metadata.getTableStatistics(session, ((TableScanNode) node.getSource()).getTable(), Constraint.alwaysTrue(), true).getRowCount(); PlanNodeStatsEstimate filteredStats = statsProvider.getStats(node); if (!filteredStats.isOutputRowCountUnknown() && !totalRowCount.isUnknown()) { diff --git a/presto-main/src/main/java/io/prestosql/sql/planner/iterative/rule/TablePushdown.java b/presto-main/src/main/java/io/prestosql/sql/planner/iterative/rule/TablePushdown.java index 6afae772a..7ebe7ac58 100644 --- a/presto-main/src/main/java/io/prestosql/sql/planner/iterative/rule/TablePushdown.java +++ b/presto-main/src/main/java/io/prestosql/sql/planner/iterative/rule/TablePushdown.java @@ -226,7 +226,7 @@ public class TablePushdown private boolean isTableWithUniqueColumns(TableScanNode tableNode) { TableHandle tableHandle = tableNode.getTable(); - TableStatistics tableStatistics = metadata.getTableStatistics(ruleContext.getSession(), tableHandle, Constraint.alwaysTrue()); + TableStatistics tableStatistics = metadata.getTableStatistics(ruleContext.getSession(), tableHandle, Constraint.alwaysTrue(), true); /* * We check here if tablestats is null or not. diff --git a/presto-main/src/main/java/io/prestosql/sql/planner/optimizations/AddReuseExchange.java b/presto-main/src/main/java/io/prestosql/sql/planner/optimizations/AddReuseExchange.java index 4daa99a45..c7809b7a7 100644 --- a/presto-main/src/main/java/io/prestosql/sql/planner/optimizations/AddReuseExchange.java +++ b/presto-main/src/main/java/io/prestosql/sql/planner/optimizations/AddReuseExchange.java @@ -216,7 +216,7 @@ public class AddReuseExchange private void visitTableScanInternal(TableScanNode node, TupleDomain newDomain) { if (!isNodeAlreadyVisited && node.getTable().getConnectorHandle().isReuseTableScanSupported()) { - TableStatistics stats = metadata.getTableStatistics(session, node.getTable(), (newDomain != null) ? new Constraint(newDomain) : Constraint.alwaysTrue()); + TableStatistics stats = metadata.getTableStatistics(session, node.getTable(), (newDomain != null) ? new Constraint(newDomain) : Constraint.alwaysTrue(), true); if (isMaxTableSizeGreaterThanSpillThreshold(node, stats)) { planNodeListHashMap.remove(WrapperScanNode.of(node)); } diff --git a/presto-main/src/main/java/io/prestosql/sql/rewrite/ShowStatsRewrite.java b/presto-main/src/main/java/io/prestosql/sql/rewrite/ShowStatsRewrite.java index f2575a2c5..54f837b8c 100644 --- a/presto-main/src/main/java/io/prestosql/sql/rewrite/ShowStatsRewrite.java +++ b/presto-main/src/main/java/io/prestosql/sql/rewrite/ShowStatsRewrite.java @@ -176,7 +176,7 @@ public class ShowStatsRewrite private Node rewriteShowStats(ShowStats node, Table table, Constraint constraint) { TableHandle tableHandle = getTableHandle(node, table.getName()); - TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, constraint); + TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, constraint, true); List statsColumnNames = buildColumnsNames(); List selectItems = buildSelectItems(statsColumnNames); TableMetadata tableMetadata = metadata.getTableMetadata(session, tableHandle); diff --git a/presto-main/src/main/java/io/prestosql/transaction/InMemoryTransactionManager.java b/presto-main/src/main/java/io/prestosql/transaction/InMemoryTransactionManager.java index ec2f3e9a3..d9c904a54 100644 --- a/presto-main/src/main/java/io/prestosql/transaction/InMemoryTransactionManager.java +++ b/presto-main/src/main/java/io/prestosql/transaction/InMemoryTransactionManager.java @@ -197,17 +197,19 @@ public class InMemoryTransactionManager } @Override - public synchronized TransactionId beginTransaction(IsolationLevel isolationLevel, boolean readOnly, boolean autoCommitContext) + public TransactionId beginTransaction(IsolationLevel isolationLevel, boolean readOnly, boolean autoCommitContext) { TransactionId transactionId = TransactionId.create(); BoundedExecutor executor = new BoundedExecutor(finishingExecutor, maxFinishingConcurrency); TransactionMetadata transactionMetadata = new TransactionMetadata(transactionId, isolationLevel, readOnly, autoCommitContext, catalogManager, executor, functionNamespaceManagers); - checkState(transactions.put(transactionId, transactionMetadata) == null, "Duplicate transaction ID: %s", transactionId); - //add transactionId to state store - if (stateStoreProvider != null && stateStoreProvider.getStateStore() != null) { - StateMap stateMap = (StateMap) stateStoreProvider.getStateStore().getStateCollection(StateStoreConstants.TRANSACTION_STATE_COLLECTION_NAME); - if (stateMap != null) { - stateMap.put(transactionId.toString(), transactionId.toString()); + synchronized (this) { + checkState(transactions.put(transactionId, transactionMetadata) == null, "Duplicate transaction ID: %s", transactionId); + //add transactionId to state store + if (stateStoreProvider != null && stateStoreProvider.getStateStore() != null) { + StateMap stateMap = (StateMap) stateStoreProvider.getStateStore().getStateCollection(StateStoreConstants.TRANSACTION_STATE_COLLECTION_NAME); + if (stateMap != null) { + stateMap.put(transactionId.toString(), transactionId.toString()); + } } } return transactionId; @@ -305,15 +307,17 @@ public class InMemoryTransactionManager tryGetTransactionMetadata(transactionId).ifPresent(TransactionMetadata::setInactive); } - private synchronized TransactionMetadata getTransactionMetadata(TransactionId transactionId) + private TransactionMetadata getTransactionMetadata(TransactionId transactionId) { TransactionMetadata transactionMetadata = transactions.get(transactionId); if (transactionMetadata == null) { // For HA use case if (stateStoreProvider != null && stateStoreProvider.getStateStore() != null) { - StateMap stateMap = (StateMap) stateStoreProvider.getStateStore().getStateCollection(StateStoreConstants.TRANSACTION_STATE_COLLECTION_NAME); - if (stateMap != null && stateMap.get(transactionId.toString()) != null) { - throw new NotInLocalTransactionException(transactionId); + synchronized (this) { + StateMap stateMap = (StateMap) stateStoreProvider.getStateStore().getStateCollection(StateStoreConstants.TRANSACTION_STATE_COLLECTION_NAME); + if (stateMap != null && stateMap.get(transactionId.toString()) != null) { + throw new NotInLocalTransactionException(transactionId); + } } } throw new NotInTransactionException(transactionId); diff --git a/presto-main/src/test/java/io/prestosql/execution/TestQueryStats.java b/presto-main/src/test/java/io/prestosql/execution/TestQueryStats.java index ffe76e006..f3bfd3d3b 100644 --- a/presto-main/src/test/java/io/prestosql/execution/TestQueryStats.java +++ b/presto-main/src/test/java/io/prestosql/execution/TestQueryStats.java @@ -166,6 +166,8 @@ public class TestQueryStats new Duration(7, NANOSECONDS), new Duration(8, NANOSECONDS), + new Duration(100, NANOSECONDS), + new Duration(100, NANOSECONDS), new Duration(100, NANOSECONDS), new Duration(200, NANOSECONDS), diff --git a/presto-main/src/test/java/io/prestosql/metadata/AbstractMockMetadata.java b/presto-main/src/test/java/io/prestosql/metadata/AbstractMockMetadata.java index fec12321c..65430d781 100644 --- a/presto-main/src/test/java/io/prestosql/metadata/AbstractMockMetadata.java +++ b/presto-main/src/test/java/io/prestosql/metadata/AbstractMockMetadata.java @@ -152,7 +152,7 @@ public abstract class AbstractMockMetadata } @Override - public TableStatistics getTableStatistics(Session session, TableHandle tableHandle, Constraint constraint) + public TableStatistics getTableStatistics(Session session, TableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { throw new UnsupportedOperationException(); } diff --git a/presto-main/src/test/java/io/prestosql/server/TestBasicQueryInfo.java b/presto-main/src/test/java/io/prestosql/server/TestBasicQueryInfo.java index c3715268b..d04aeb684 100644 --- a/presto-main/src/test/java/io/prestosql/server/TestBasicQueryInfo.java +++ b/presto-main/src/test/java/io/prestosql/server/TestBasicQueryInfo.java @@ -66,6 +66,8 @@ public class TestBasicQueryInfo Duration.valueOf("10m"), Duration.valueOf("11m"), Duration.valueOf("12m"), + Duration.valueOf("12m"), + Duration.valueOf("12m"), 13, 14, 15, diff --git a/presto-main/src/test/java/io/prestosql/server/TestQueryStateInfo.java b/presto-main/src/test/java/io/prestosql/server/TestQueryStateInfo.java index 0a16879c6..79f7f30c6 100644 --- a/presto-main/src/test/java/io/prestosql/server/TestQueryStateInfo.java +++ b/presto-main/src/test/java/io/prestosql/server/TestQueryStateInfo.java @@ -116,6 +116,8 @@ public class TestQueryStateInfo Duration.valueOf("9m"), Duration.valueOf("10m"), Duration.valueOf("11m"), + Duration.valueOf("11m"), + Duration.valueOf("11m"), Duration.valueOf("12m"), 13, 14, diff --git a/presto-product-tests/src/test/resources/io/prestosql/tests/querystats/single_query_info_response.json b/presto-product-tests/src/test/resources/io/prestosql/tests/querystats/single_query_info_response.json index 6195d1df0..816d155f1 100644 --- a/presto-product-tests/src/test/resources/io/prestosql/tests/querystats/single_query_info_response.json +++ b/presto-product-tests/src/test/resources/io/prestosql/tests/querystats/single_query_info_response.json @@ -34,6 +34,8 @@ "analysisTime": "7.47ms", "distributedPlanningTime": "311.77us", "totalPlanningTime": "9.99ms", + "totalLogicalPlanningTime": "3.33ms", + "totalSyntaxAnalysisTime": "1.11ms", "finishingTime": "17.00ms", "totalTasks": 1, "runningTasks": 0, diff --git a/presto-spi/src/main/java/io/prestosql/spi/connector/CachedConnectorMetadata.java b/presto-spi/src/main/java/io/prestosql/spi/connector/CachedConnectorMetadata.java index 5300e50e2..a1a6c4397 100644 --- a/presto-spi/src/main/java/io/prestosql/spi/connector/CachedConnectorMetadata.java +++ b/presto-spi/src/main/java/io/prestosql/spi/connector/CachedConnectorMetadata.java @@ -329,18 +329,18 @@ public class CachedConnectorMetadata @Override public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, - Constraint constraint) + Constraint constraint, boolean includeColumnStatistics) { Optional cacheOpt = getOrCreateCache(session); if (!cacheOpt.isPresent()) { - return logAndDelegate("getTableStatistics", () -> delegate.getTableStatistics(session, tableHandle, constraint)); + return logAndDelegate("getTableStatistics", () -> delegate.getTableStatistics(session, tableHandle, constraint, includeColumnStatistics)); } try { return cacheOpt.get().getTableStatistics().get(tableHandle.getSchemaPrefixedTableName(), () -> { TableStatistics tableStatistics = - logAndDelegate("getTableStatistics", () -> delegate.getTableStatistics(session, tableHandle, constraint)); + logAndDelegate("getTableStatistics", () -> delegate.getTableStatistics(session, tableHandle, constraint, includeColumnStatistics)); if (tableStatistics == null) { throw new Exception(); @@ -350,7 +350,7 @@ public class CachedConnectorMetadata }); } catch (Exception e) { - return logAndDelegate("getTableStatistics", () -> delegate.getTableStatistics(session, tableHandle, constraint)); + return logAndDelegate("getTableStatistics", () -> delegate.getTableStatistics(session, tableHandle, constraint, includeColumnStatistics)); } } diff --git a/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorMetadata.java b/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorMetadata.java index 952b28bc2..db253e016 100644 --- a/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorMetadata.java +++ b/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorMetadata.java @@ -224,7 +224,7 @@ public interface ConnectorMetadata /** * Get statistics for table for given filtering constraint. */ - default TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint) + default TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { return TableStatistics.empty(); } diff --git a/presto-spi/src/main/java/io/prestosql/spi/connector/classloader/ClassLoaderSafeConnectorMetadata.java b/presto-spi/src/main/java/io/prestosql/spi/connector/classloader/ClassLoaderSafeConnectorMetadata.java index 09fba3195..d5e46227a 100644 --- a/presto-spi/src/main/java/io/prestosql/spi/connector/classloader/ClassLoaderSafeConnectorMetadata.java +++ b/presto-spi/src/main/java/io/prestosql/spi/connector/classloader/ClassLoaderSafeConnectorMetadata.java @@ -272,10 +272,10 @@ public class ClassLoaderSafeConnectorMetadata } @Override - public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint) + public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { try (ThreadContextClassLoader ignored = new ThreadContextClassLoader(classLoader)) { - return delegate.getTableStatistics(session, tableHandle, constraint); + return delegate.getTableStatistics(session, tableHandle, constraint, includeColumnStatistics); } } diff --git a/presto-spi/src/main/java/io/prestosql/spi/statistics/TableStatistics.java b/presto-spi/src/main/java/io/prestosql/spi/statistics/TableStatistics.java index 33d262c0d..c5814ae71 100644 --- a/presto-spi/src/main/java/io/prestosql/spi/statistics/TableStatistics.java +++ b/presto-spi/src/main/java/io/prestosql/spi/statistics/TableStatistics.java @@ -29,6 +29,8 @@ public final class TableStatistics private static final TableStatistics EMPTY = TableStatistics.builder().build(); private final Estimate rowCount; + private final long fileCount; + private final long onDiskDataSizeInBytes; private final Map columnStatistics; public static TableStatistics empty() @@ -36,9 +38,12 @@ public final class TableStatistics return EMPTY; } - public TableStatistics(Estimate rowCount, Map columnStatistics) + // added parameters fileCount and onDiskDataSizeInBytes used as check for invalidating tableStatisticsCache. + public TableStatistics(Estimate rowCount, long fileCount, long onDiskDataSizeInBytes, Map columnStatistics) { this.rowCount = requireNonNull(rowCount, "rowCount can not be null"); + this.fileCount = requireNonNull(fileCount, "fileCount can not be null"); + this.onDiskDataSizeInBytes = requireNonNull(onDiskDataSizeInBytes, "onDiskDataSizeInBytes can not be null"); if (!rowCount.isUnknown() && rowCount.getValue() < 0) { throw new IllegalArgumentException(format("rowCount must be greater than or equal to 0: %s", rowCount.getValue())); } @@ -50,6 +55,16 @@ public final class TableStatistics return rowCount; } + public long getFileCount() + { + return fileCount; + } + + public long getOnDiskDataSizeInBytes() + { + return onDiskDataSizeInBytes; + } + public Map getColumnStatistics() { return columnStatistics; @@ -92,6 +107,8 @@ public final class TableStatistics public static final class Builder { private Estimate rowCount = Estimate.unknown(); + private long fileCount; + private long onDiskDataSizeInBytes; private Map columnStatisticsMap = new LinkedHashMap<>(); public Builder setRowCount(Estimate rowCount) @@ -100,6 +117,18 @@ public final class TableStatistics return this; } + public Builder setFileCount(long fileCount) + { + this.fileCount = requireNonNull(fileCount, "fileCount can not be null"); + return this; + } + + public Builder setOnDiskDataSizeInBytes(long onDiskDataSizeInBytes) + { + this.onDiskDataSizeInBytes = requireNonNull(onDiskDataSizeInBytes, "onDiskDataSizeInBytes can not be null"); + return this; + } + public Builder setColumnStatistics(ColumnHandle columnHandle, ColumnStatistics columnStatistics) { requireNonNull(columnHandle, "columnHandle can not be null"); @@ -110,7 +139,7 @@ public final class TableStatistics public TableStatistics build() { - return new TableStatistics(rowCount, columnStatisticsMap); + return new TableStatistics(rowCount, fileCount, onDiskDataSizeInBytes, columnStatisticsMap); } } } diff --git a/presto-tpcds/src/main/java/io/prestosql/plugin/tpcds/TpcdsMetadata.java b/presto-tpcds/src/main/java/io/prestosql/plugin/tpcds/TpcdsMetadata.java index bc892888f..10956c862 100644 --- a/presto-tpcds/src/main/java/io/prestosql/plugin/tpcds/TpcdsMetadata.java +++ b/presto-tpcds/src/main/java/io/prestosql/plugin/tpcds/TpcdsMetadata.java @@ -154,7 +154,7 @@ public class TpcdsMetadata } @Override - public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint) + public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { TpcdsTableHandle tpcdsTableHandle = (TpcdsTableHandle) tableHandle; diff --git a/presto-tpcds/src/test/java/io/prestosql/plugin/tpcds/TestTpcdsMetadataStatistics.java b/presto-tpcds/src/test/java/io/prestosql/plugin/tpcds/TestTpcdsMetadataStatistics.java index 96c703e34..1f52a4857 100644 --- a/presto-tpcds/src/test/java/io/prestosql/plugin/tpcds/TestTpcdsMetadataStatistics.java +++ b/presto-tpcds/src/test/java/io/prestosql/plugin/tpcds/TestTpcdsMetadataStatistics.java @@ -50,7 +50,7 @@ public class TestTpcdsMetadataStatistics .forEach(table -> { SchemaTableName schemaTableName = new SchemaTableName(schemaName, table.getName()); ConnectorTableHandle tableHandle = metadata.getTableHandle(session, schemaTableName); - TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, alwaysTrue()); + TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, alwaysTrue(), true); assertTrue(tableStatistics.getRowCount().isUnknown()); assertTrue(tableStatistics.getColumnStatistics().isEmpty()); })); @@ -64,7 +64,7 @@ public class TestTpcdsMetadataStatistics .forEach(table -> { SchemaTableName schemaTableName = new SchemaTableName(schemaName, table.getName()); ConnectorTableHandle tableHandle = metadata.getTableHandle(session, schemaTableName); - TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, alwaysTrue()); + TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, alwaysTrue(), true); assertFalse(tableStatistics.getRowCount().isUnknown()); for (ColumnHandle column : metadata.getColumnHandles(session, tableHandle).values()) { assertTrue(tableStatistics.getColumnStatistics().containsKey(column)); @@ -78,7 +78,7 @@ public class TestTpcdsMetadataStatistics { SchemaTableName schemaTableName = new SchemaTableName("sf1", Table.CALL_CENTER.getName()); ConnectorTableHandle tableHandle = metadata.getTableHandle(session, schemaTableName); - TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, alwaysTrue()); + TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, alwaysTrue(), true); estimateAssertion.assertClose(tableStatistics.getRowCount(), Estimate.of(6), "Row count does not match"); @@ -148,7 +148,7 @@ public class TestTpcdsMetadataStatistics { SchemaTableName schemaTableName = new SchemaTableName("sf1", Table.WEB_SITE.getName()); ConnectorTableHandle tableHandle = metadata.getTableHandle(session, schemaTableName); - TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, alwaysTrue()); + TableStatistics tableStatistics = metadata.getTableStatistics(session, tableHandle, alwaysTrue(), true); Map columnHandles = metadata.getColumnHandles(session, tableHandle); diff --git a/presto-tpch/src/main/java/io/prestosql/plugin/tpch/TpchMetadata.java b/presto-tpch/src/main/java/io/prestosql/plugin/tpch/TpchMetadata.java index 4bd2def7c..bdc738eda 100644 --- a/presto-tpch/src/main/java/io/prestosql/plugin/tpch/TpchMetadata.java +++ b/presto-tpch/src/main/java/io/prestosql/plugin/tpch/TpchMetadata.java @@ -257,7 +257,7 @@ public class TpchMetadata } @Override - public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint) + public TableStatistics getTableStatistics(ConnectorSession session, ConnectorTableHandle tableHandle, Constraint constraint, boolean includeColumnStatistics) { TpchTableHandle tpchTableHandle = (TpchTableHandle) tableHandle; String tableName = tpchTableHandle.getTableName(); diff --git a/presto-tpch/src/test/java/io/prestosql/plugin/tpch/TestTpchMetadata.java b/presto-tpch/src/test/java/io/prestosql/plugin/tpch/TestTpchMetadata.java index a67fc8f4d..0a53a3146 100644 --- a/presto-tpch/src/test/java/io/prestosql/plugin/tpch/TestTpchMetadata.java +++ b/presto-tpch/src/test/java/io/prestosql/plugin/tpch/TestTpchMetadata.java @@ -187,7 +187,7 @@ public class TestTpchMetadata private void testTableStats(String schema, TpchTable table, Constraint constraint, double expectedRowCount) { TpchTableHandle tableHandle = tpchMetadata.getTableHandle(session, new SchemaTableName(schema, table.getTableName())); - TableStatistics tableStatistics = tpchMetadata.getTableStatistics(session, tableHandle, constraint); + TableStatistics tableStatistics = tpchMetadata.getTableStatistics(session, tableHandle, constraint, true); double actualRowCountValue = tableStatistics.getRowCount().getValue(); assertEquals(tableStatistics.getRowCount(), Estimate.of(actualRowCountValue)); @@ -197,7 +197,7 @@ public class TestTpchMetadata private void testNoTableStats(String schema, TpchTable table) { TpchTableHandle tableHandle = tpchMetadata.getTableHandle(session, new SchemaTableName(schema, table.getTableName())); - TableStatistics tableStatistics = tpchMetadata.getTableStatistics(session, tableHandle, alwaysTrue()); + TableStatistics tableStatistics = tpchMetadata.getTableStatistics(session, tableHandle, alwaysTrue(), true); assertTrue(tableStatistics.getRowCount().isUnknown()); } @@ -288,7 +288,7 @@ public class TestTpchMetadata private void testColumnStats(String schema, TpchTable table, TpchColumn column, Constraint constraint, ColumnStatistics expected) { TpchTableHandle tableHandle = tpchMetadata.getTableHandle(session, new SchemaTableName(schema, table.getTableName())); - TableStatistics tableStatistics = tpchMetadata.getTableStatistics(session, tableHandle, constraint); + TableStatistics tableStatistics = tpchMetadata.getTableStatistics(session, tableHandle, constraint, true); ColumnHandle columnHandle = tpchMetadata.getColumnHandles(session, tableHandle).get(column.getSimplifiedColumnName()); ColumnStatistics actual = tableStatistics.getColumnStatistics().get(columnHandle);