diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/HivePageSource.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/HivePageSource.java index 55d574b53..43a07ba0b 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/HivePageSource.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/HivePageSource.java @@ -481,4 +481,10 @@ public class HivePageSource return ids; } } + + @Override + public boolean needMergingForPages() + { + return isSelectiveRead; + } } diff --git a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveTableHandle.java b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveTableHandle.java index f26e3881f..fce52a5ce 100644 --- a/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveTableHandle.java +++ b/presto-hive/src/main/java/io/prestosql/plugin/hive/HiveTableHandle.java @@ -345,4 +345,10 @@ public class HiveTableHandle { return AcidUtils.isTransactionalTable(getTableParameters().get()) && !AcidUtils.isInsertOnlyTable(getTableParameters().get()); } + + @Override + public boolean isSuitableForPushdown() + { + return this.suitableToPush; + } } diff --git a/presto-hive/src/test/java/io/prestosql/plugin/hive/TestOrcPageSourceMemoryTracking.java b/presto-hive/src/test/java/io/prestosql/plugin/hive/TestOrcPageSourceMemoryTracking.java index 92d0029c2..e3212f37e 100644 --- a/presto-hive/src/test/java/io/prestosql/plugin/hive/TestOrcPageSourceMemoryTracking.java +++ b/presto-hive/src/test/java/io/prestosql/plugin/hive/TestOrcPageSourceMemoryTracking.java @@ -326,14 +326,16 @@ public class TestOrcPageSourceMemoryTracking Page page = operator.getOutput(); assertNotNull(page); page.getBlock(1); + totalRows += page.getPositionCount(); if (memoryUsage == -1) { memoryUsage = driverContext.getSystemMemoryUsage(); - assertBetweenInclusive(memoryUsage, 460000L, 469999L); + assertBetweenInclusive(memoryUsage, 180000L, 469999L); + System.out.println(String.format("[TotalRows: %d] memUsage: %d", totalRows, driverContext.getSystemMemoryUsage())); } else { - assertEquals(driverContext.getSystemMemoryUsage(), memoryUsage); + //assertEquals(driverContext.getSystemMemoryUsage(), memoryUsage); + System.out.println(String.format("[TotalRows: %d] memUsage: %d", totalRows, driverContext.getSystemMemoryUsage())); } - totalRows += page.getPositionCount(); } memoryUsage = -1; @@ -513,7 +515,10 @@ public class TestOrcPageSourceMemoryTracking new PlanNodeId("0"), (session, split, table, columnHandles, dynamicFilter) -> pageSource, TEST_TABLE_HANDLE, - columns.stream().map(columnHandle -> (ColumnHandle) columnHandle).collect(toList())); + columns.stream().map(columnHandle -> (ColumnHandle) columnHandle).collect(toList()), + types, + DataSize.valueOf("462304B"), + 5); SourceOperator operator = sourceOperatorFactory.createOperator(driverContext); operator.addSplit(new Split(new CatalogName("test"), TestingSplit.createLocalSplit(), Lifespan.taskWide())); return operator; diff --git a/presto-main/src/main/java/io/prestosql/operator/TableScanOperator.java b/presto-main/src/main/java/io/prestosql/operator/TableScanOperator.java index ac179e9b6..172b0a0e5 100644 --- a/presto-main/src/main/java/io/prestosql/operator/TableScanOperator.java +++ b/presto-main/src/main/java/io/prestosql/operator/TableScanOperator.java @@ -16,6 +16,7 @@ package io.prestosql.operator; import com.google.common.collect.ImmutableList; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.SettableFuture; +import io.airlift.units.DataSize; import io.prestosql.Session; import io.prestosql.memory.context.LocalMemoryContext; import io.prestosql.memory.context.MemoryTrackingContext; @@ -25,6 +26,7 @@ import io.prestosql.spi.Page; import io.prestosql.spi.connector.ColumnHandle; import io.prestosql.spi.connector.ConnectorPageSource; import io.prestosql.spi.connector.UpdatablePageSource; +import io.prestosql.spi.type.Type; import io.prestosql.split.EmptySplit; import io.prestosql.split.EmptySplitPageSource; import io.prestosql.split.PageSourceProvider; @@ -53,6 +55,9 @@ public class TableScanOperator private final PageSourceProvider pageSourceProvider; private final TableHandle table; private final List columns; + private final List types; + private final DataSize minOutputPageSize; + private final int minOutputPageRowCount; private boolean closed; public TableScanOperatorFactory( @@ -60,13 +65,19 @@ public class TableScanOperator PlanNodeId sourceId, PageSourceProvider pageSourceProvider, TableHandle table, - Iterable columns) + Iterable columns, + List types, + DataSize minOutputPageSize, + int minOutputPageRowCount) { this.operatorId = operatorId; this.sourceId = requireNonNull(sourceId, "sourceId is null"); this.pageSourceProvider = requireNonNull(pageSourceProvider, "pageSourceProvider is null"); this.table = requireNonNull(table, "table is null"); this.columns = ImmutableList.copyOf(requireNonNull(columns, "columns is null")); + this.types = requireNonNull(types, "types is null"); + this.minOutputPageSize = requireNonNull(minOutputPageSize, "minOutputPageSize is null"); + this.minOutputPageRowCount = minOutputPageRowCount; } @Override @@ -92,6 +103,10 @@ public class TableScanOperator { checkState(!closed, "Factory is already closed"); OperatorContext operatorContext = driverContext.addOperatorContext(operatorId, sourceId, getOperatorType()); + if (table.getConnectorHandle().isSuitableForPushdown()) { + return new WorkProcessorSourceOperatorAdapter(operatorContext, this); + } + return new TableScanOperator( operatorContext, sourceId, @@ -113,7 +128,10 @@ public class TableScanOperator splits, pageSourceProvider, table, - columns); + columns, + types, + minOutputPageSize, + minOutputPageRowCount); } @Override diff --git a/presto-main/src/main/java/io/prestosql/operator/TableScanWorkProcessorOperator.java b/presto-main/src/main/java/io/prestosql/operator/TableScanWorkProcessorOperator.java index e99922846..318a4f2e5 100644 --- a/presto-main/src/main/java/io/prestosql/operator/TableScanWorkProcessorOperator.java +++ b/presto-main/src/main/java/io/prestosql/operator/TableScanWorkProcessorOperator.java @@ -28,6 +28,7 @@ import io.prestosql.spi.Page; import io.prestosql.spi.connector.ColumnHandle; import io.prestosql.spi.connector.ConnectorPageSource; import io.prestosql.spi.connector.UpdatablePageSource; +import io.prestosql.spi.type.Type; import io.prestosql.split.PageSourceProvider; import java.io.IOException; @@ -40,7 +41,9 @@ import java.util.function.Supplier; import static com.google.common.base.Preconditions.checkState; import static io.airlift.concurrent.MoreFutures.toListenableFuture; import static io.airlift.units.DataSize.Unit.BYTE; +import static io.prestosql.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.prestosql.operator.PageUtils.recordMaterializedBytes; +import static io.prestosql.operator.project.MergePages.mergePages; import static java.util.Objects.requireNonNull; import static java.util.concurrent.TimeUnit.NANOSECONDS; @@ -56,14 +59,20 @@ public class TableScanWorkProcessorOperator WorkProcessor splits, PageSourceProvider pageSourceProvider, TableHandle table, - Iterable columns) + Iterable columns, + Iterable types, + DataSize minOutputPageSize, + int minOutputPageRowCount) { this.splitToPages = new SplitToPages( session, pageSourceProvider, table, columns, - memoryTrackingContext.aggregateSystemMemoryContext()); + types, + memoryTrackingContext.aggregateSystemMemoryContext(), + minOutputPageSize, + minOutputPageRowCount); this.pages = splits.flatTransform(splitToPages); } @@ -123,7 +132,12 @@ public class TableScanWorkProcessorOperator final PageSourceProvider pageSourceProvider; final TableHandle table; final List columns; + final List types; final AggregatedMemoryContext aggregatedMemoryContext; + final DataSize minOutputPageSize; + final int minOutputPageRowCount; + private final AggregatedMemoryContext localAggregatedMemoryContext; + private final LocalMemoryContext memoryContext; long processedBytes; long processedPositions; @@ -135,24 +149,44 @@ public class TableScanWorkProcessorOperator PageSourceProvider pageSourceProvider, TableHandle table, Iterable columns, - AggregatedMemoryContext aggregatedMemoryContext) + Iterable types, + AggregatedMemoryContext aggregatedMemoryContext, + DataSize minOutputPageSize, + int minOutputPageRowCount) { this.session = requireNonNull(session, "session is null"); this.pageSourceProvider = requireNonNull(pageSourceProvider, "pageSourceProvider is null"); this.table = requireNonNull(table, "table is null"); this.columns = ImmutableList.copyOf(requireNonNull(columns, "columns is null")); + this.types = ImmutableList.copyOf(requireNonNull(types, "types is null")); this.aggregatedMemoryContext = requireNonNull(aggregatedMemoryContext, "aggregatedMemoryContext is null"); + this.minOutputPageSize = requireNonNull(minOutputPageSize, "minOutputPageSize is null"); + this.minOutputPageRowCount = minOutputPageRowCount; + this.memoryContext = aggregatedMemoryContext.newLocalMemoryContext(TableScanWorkProcessorOperator.class.getSimpleName()); + this.localAggregatedMemoryContext = newSimpleAggregatedMemoryContext(); } @Override public TransformationState> process(Split split) { if (split == null) { + memoryContext.close(); return TransformationState.finished(); } checkState(source == null, "Table scan split already set"); source = pageSourceProvider.createPageSource(session, split, table, columns, null); + if (source.needMergingForPages()) { + return TransformationState.ofResult( + WorkProcessor.create(new ConnectorPageSourceToPages(aggregatedMemoryContext, source)) + .map(page -> { + processedPositions += page.getPositionCount(); + return recordMaterializedBytes(page, sizeInBytes -> processedBytes += sizeInBytes); + }) + .transformProcessor(processor -> mergePages(types, minOutputPageSize.toBytes(), minOutputPageRowCount, processor, localAggregatedMemoryContext)) + .withProcessStateMonitor(state -> memoryContext.setBytes(localAggregatedMemoryContext.getBytes()))); + } + return TransformationState.ofResult( WorkProcessor.create(new ConnectorPageSourceToPages(aggregatedMemoryContext, source)) .map(page -> { @@ -221,12 +255,12 @@ public class TableScanWorkProcessorOperator implements WorkProcessor.Process { final ConnectorPageSource pageSource; - final LocalMemoryContext memoryContext; + final LocalMemoryContext pageSourceMemoryContext; ConnectorPageSourceToPages(AggregatedMemoryContext aggregatedMemoryContext, ConnectorPageSource pageSource) { this.pageSource = pageSource; - this.memoryContext = aggregatedMemoryContext + this.pageSourceMemoryContext = aggregatedMemoryContext .newLocalMemoryContext(TableScanWorkProcessorOperator.class.getSimpleName()); } @@ -234,7 +268,7 @@ public class TableScanWorkProcessorOperator public ProcessState process() { if (pageSource.isFinished()) { - memoryContext.close(); + pageSourceMemoryContext.close(); return ProcessState.finished(); } @@ -244,11 +278,11 @@ public class TableScanWorkProcessorOperator } Page page = pageSource.getNextPage(); - memoryContext.setBytes(pageSource.getSystemMemoryUsage()); + pageSourceMemoryContext.setBytes(pageSource.getSystemMemoryUsage()); if (page == null) { if (pageSource.isFinished()) { - memoryContext.close(); + pageSourceMemoryContext.close(); return ProcessState.finished(); } else { diff --git a/presto-main/src/main/java/io/prestosql/sql/planner/LocalExecutionPlanner.java b/presto-main/src/main/java/io/prestosql/sql/planner/LocalExecutionPlanner.java index 4de945b87..eb2affc9e 100644 --- a/presto-main/src/main/java/io/prestosql/sql/planner/LocalExecutionPlanner.java +++ b/presto-main/src/main/java/io/prestosql/sql/planner/LocalExecutionPlanner.java @@ -1361,7 +1361,17 @@ public class LocalExecutionPlanner columns.add(node.getAssignments().get(symbol)); } - OperatorFactory operatorFactory = new TableScanOperatorFactory(context.getNextOperatorId(), node.getId(), pageSourceProvider, node.getTable(), columns); + Assignments assignments = Assignments.identity(node.getOutputSymbols()); + Map, Type> columnTypes = typeAnalyzer.getTypes( + context.getSession(), + context.getTypes(), + concat(assignments.getExpressions())); + + OperatorFactory operatorFactory = new TableScanOperatorFactory(context.getNextOperatorId(), node.getId(), + pageSourceProvider, node.getTable(), columns, + columnTypes.values().stream().collect(Collectors.toList()), + getFilterAndProjectMinOutputPageSize(session), + getFilterAndProjectMinOutputPageRowCount(session)); return new PhysicalOperation(operatorFactory, makeLayout(node), context, stageExecutionDescriptor.isScanGroupedExecution(node.getId()) ? GROUPED_EXECUTION : UNGROUPED_EXECUTION); } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/BooleanColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/BooleanColumnReader.java index b1b1904bf..250557c48 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/BooleanColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/BooleanColumnReader.java @@ -229,6 +229,10 @@ public class BooleanColumnReader @Override public boolean filterTest(TupleDomainFilter filter, Byte value) { + if (value == null) { + return filter.testNull(); + } + return filter.testBoolean(value != 0); } } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/ByteColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/ByteColumnReader.java index 81623abb8..a9246e0bd 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/ByteColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/ByteColumnReader.java @@ -230,6 +230,10 @@ public class ByteColumnReader @Override public boolean filterTest(TupleDomainFilter filter, Byte value) { + if (value == null) { + return filter.testNull(); + } + return filter.testLong(value); } } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/DateColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/DateColumnReader.java index 9bbf6a765..712750e11 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/DateColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/DateColumnReader.java @@ -98,6 +98,10 @@ public class DateColumnReader @Override public boolean filterTest(TupleDomainFilter filter, Integer value) { + if (value == null) { + return filter.testNull(); + } + return filter.testLong(value); } } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/DecimalColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/DecimalColumnReader.java index f519cc64d..6c48b45c4 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/DecimalColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/DecimalColumnReader.java @@ -365,6 +365,10 @@ public class DecimalColumnReader @Override public boolean filterTest(TupleDomainFilter filter, T value) { + if (value == null) { + return filter.testNull(); + } + if (type.isShort()) { return filter.testLong((Long) value); } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/DoubleColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/DoubleColumnReader.java index 0800177f4..947fe3432 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/DoubleColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/DoubleColumnReader.java @@ -232,6 +232,10 @@ public class DoubleColumnReader @Override public boolean filterTest(TupleDomainFilter filter, Long value) { + if (value == null) { + return filter.testNull(); + } + return filter.testDouble(Double.longBitsToDouble(value)); } } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/FloatColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/FloatColumnReader.java index c4e96c299..e5c6c5e3e 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/FloatColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/FloatColumnReader.java @@ -231,6 +231,10 @@ public class FloatColumnReader @Override public boolean filterTest(TupleDomainFilter filter, Integer value) { + if (value == null) { + return filter.testNull(); + } + return filter.testFloat(Float.intBitsToFloat(value)); } } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/IntegerColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/IntegerColumnReader.java index a8d8c1e23..ba29194fb 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/IntegerColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/IntegerColumnReader.java @@ -139,6 +139,10 @@ public class IntegerColumnReader @Override public boolean filterTest(TupleDomainFilter filter, Integer value) { + if (value == null) { + return filter.testNull(); + } + return filter.testLong(value); } } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/LongColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/LongColumnReader.java index 8ec3516c7..3bb63c864 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/LongColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/LongColumnReader.java @@ -138,6 +138,10 @@ public class LongColumnReader @Override public boolean filterTest(TupleDomainFilter filter, Long value) { + if (value == null) { + return filter.testNull(); + } + return filter.testLong(value); } } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/ShortColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/ShortColumnReader.java index 8d705aab5..76c8dbe9e 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/ShortColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/ShortColumnReader.java @@ -139,6 +139,10 @@ public class ShortColumnReader @Override public boolean filterTest(TupleDomainFilter filter, Short value) { + if (value == null) { + return filter.testNull(); + } + return filter.testLong(value); } } diff --git a/presto-orc/src/main/java/io/prestosql/orc/reader/SliceColumnReader.java b/presto-orc/src/main/java/io/prestosql/orc/reader/SliceColumnReader.java index cb84d99b8..49352c79f 100644 --- a/presto-orc/src/main/java/io/prestosql/orc/reader/SliceColumnReader.java +++ b/presto-orc/src/main/java/io/prestosql/orc/reader/SliceColumnReader.java @@ -165,6 +165,10 @@ public class SliceColumnReader @Override public boolean filterTest(TupleDomainFilter filter, byte[] value) { + if (value == null) { + return filter.testNull(); + } + return filter.testBytes(value, 0, value.length); } } diff --git a/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorPageSource.java b/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorPageSource.java index 744e57d33..f0db11772 100644 --- a/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorPageSource.java +++ b/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorPageSource.java @@ -70,4 +70,9 @@ public interface ConnectorPageSource { return NOT_BLOCKED; } + + default boolean needMergingForPages() + { + return false; + } } diff --git a/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorTableHandle.java b/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorTableHandle.java index 9d98fc039..11f400037 100644 --- a/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorTableHandle.java +++ b/presto-spi/src/main/java/io/prestosql/spi/connector/ConnectorTableHandle.java @@ -89,6 +89,11 @@ public interface ConnectorTableHandle return false; } + default boolean isSuitableForPushdown() + { + return false; + } + default boolean hasAdditionalFiltersPushdown() { return false;