diff --git a/core/trino-main/src/main/java/io/trino/client/direct/DirectTrinoClient.java b/core/trino-main/src/main/java/io/trino/client/direct/DirectTrinoClient.java index 6f94af205ae1..68192e3f2154 100644 --- a/core/trino-main/src/main/java/io/trino/client/direct/DirectTrinoClient.java +++ b/core/trino-main/src/main/java/io/trino/client/direct/DirectTrinoClient.java @@ -23,7 +23,6 @@ import io.trino.execution.QueryManager; import io.trino.execution.QueryState; import io.trino.execution.buffer.PageDeserializer; -import io.trino.execution.buffer.PagesSerdeFactory; import io.trino.memory.context.SimpleLocalMemoryContext; import io.trino.operator.DirectExchangeClient; import io.trino.operator.DirectExchangeClientSupplier; @@ -43,10 +42,10 @@ import java.util.concurrent.ExecutionException; import static io.airlift.concurrent.MoreFutures.whenAnyComplete; -import static io.trino.SystemSessionProperties.getExchangeCompressionCodec; import static io.trino.SystemSessionProperties.getRetryPolicy; import static io.trino.execution.QueryState.FAILED; import static io.trino.execution.QueryState.FINISHING; +import static io.trino.execution.buffer.PagesSerdes.createExchangePagesSerdeFactory; import static io.trino.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.trino.spi.StandardErrorCode.GENERIC_INTERNAL_ERROR; import static java.util.Objects.requireNonNull; @@ -94,7 +93,7 @@ public DispatchQuery execute(SessionContext sessionContext, @Language("SQL") Str } }); - PageDeserializer pageDeserializer = new PagesSerdeFactory(blockEncodingSerde, getExchangeCompressionCodec(dispatchQuery.getSession())).createDeserializer(Optional.empty()); + PageDeserializer pageDeserializer = createExchangePagesSerdeFactory(blockEncodingSerde, dispatchQuery.getSession()).createDeserializer(Optional.empty()); for (QueryState state = queryManager.getQueryState(queryId); (state != FAILED) && !exchangeClient.isFinished() && diff --git a/core/trino-main/src/main/java/io/trino/execution/buffer/PagesSerdeFactory.java b/core/trino-main/src/main/java/io/trino/execution/buffer/PagesSerdeFactory.java index c0bdd1e5fb63..0f6c35815980 100644 --- a/core/trino-main/src/main/java/io/trino/execution/buffer/PagesSerdeFactory.java +++ b/core/trino-main/src/main/java/io/trino/execution/buffer/PagesSerdeFactory.java @@ -39,7 +39,8 @@ public class PagesSerdeFactory private final CompressionCodec compressionCodec; private final int blockSizeInBytes; - public PagesSerdeFactory(BlockEncodingSerde blockEncodingSerde, CompressionCodec compressionCodec) + // created via PagesSerdes.create* + PagesSerdeFactory(BlockEncodingSerde blockEncodingSerde, CompressionCodec compressionCodec) { this(blockEncodingSerde, compressionCodec, SERIALIZED_PAGE_DEFAULT_BLOCK_SIZE_IN_BYTES); } diff --git a/core/trino-main/src/main/java/io/trino/execution/buffer/PagesSerdes.java b/core/trino-main/src/main/java/io/trino/execution/buffer/PagesSerdes.java new file mode 100644 index 000000000000..cbd463b2d82f --- /dev/null +++ b/core/trino-main/src/main/java/io/trino/execution/buffer/PagesSerdes.java @@ -0,0 +1,34 @@ +/* + * 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.trino.execution.buffer; + +import io.trino.Session; +import io.trino.spi.block.BlockEncodingSerde; + +import static io.trino.SystemSessionProperties.getExchangeCompressionCodec; + +public final class PagesSerdes +{ + private PagesSerdes() {} + + public static PagesSerdeFactory createExchangePagesSerdeFactory(BlockEncodingSerde blockEncodingSerde, Session session) + { + return new PagesSerdeFactory(blockEncodingSerde, getExchangeCompressionCodec(session)); + } + + public static PagesSerdeFactory createSpillingPagesSerdeFactory(BlockEncodingSerde blockEncodingSerde, CompressionCodec compressionCodec) + { + return new PagesSerdeFactory(blockEncodingSerde, compressionCodec); + } +} diff --git a/core/trino-main/src/main/java/io/trino/server/protocol/Query.java b/core/trino-main/src/main/java/io/trino/server/protocol/Query.java index abc67b3cf034..a6cd9b39f09a 100644 --- a/core/trino-main/src/main/java/io/trino/server/protocol/Query.java +++ b/core/trino-main/src/main/java/io/trino/server/protocol/Query.java @@ -42,7 +42,6 @@ import io.trino.execution.QueryState; import io.trino.execution.StageId; import io.trino.execution.buffer.PageDeserializer; -import io.trino.execution.buffer.PagesSerdeFactory; import io.trino.memory.context.SimpleLocalMemoryContext; import io.trino.operator.DirectExchangeClientSupplier; import io.trino.server.ExternalUriInfo; @@ -86,10 +85,10 @@ import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static io.airlift.concurrent.MoreFutures.addTimeout; import static io.airlift.units.DataSize.Unit.MEGABYTE; -import static io.trino.SystemSessionProperties.getExchangeCompressionCodec; import static io.trino.SystemSessionProperties.getRetryPolicy; import static io.trino.execution.QueryState.FAILED; import static io.trino.execution.QueryState.FINISHING; +import static io.trino.execution.buffer.PagesSerdes.createExchangePagesSerdeFactory; import static io.trino.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.trino.server.protocol.ProtocolUtil.createColumn; import static io.trino.server.protocol.ProtocolUtil.toStatementStats; @@ -262,7 +261,7 @@ private Query( this.resultsProcessorExecutor = resultsProcessorExecutor; this.timeoutExecutor = timeoutExecutor; this.supportsParametricDateTime = session.getClientCapabilities().contains(ClientCapabilities.PARAMETRIC_DATETIME.toString()); - deserializer = new PagesSerdeFactory(blockEncodingSerde, getExchangeCompressionCodec(session)) + deserializer = createExchangePagesSerdeFactory(blockEncodingSerde, session) .createDeserializer(session.getExchangeEncryptionKey().map(Ciphers::deserializeAesEncryptionKey)); } diff --git a/core/trino-main/src/main/java/io/trino/spiller/FileSingleStreamSpillerFactory.java b/core/trino-main/src/main/java/io/trino/spiller/FileSingleStreamSpillerFactory.java index eed4a9426318..75eb0651dea4 100644 --- a/core/trino-main/src/main/java/io/trino/spiller/FileSingleStreamSpillerFactory.java +++ b/core/trino-main/src/main/java/io/trino/spiller/FileSingleStreamSpillerFactory.java @@ -46,6 +46,7 @@ import static io.airlift.concurrent.Threads.daemonThreadsNamed; import static io.trino.FeaturesConfig.SPILLER_SPILL_PATH; import static io.trino.cache.SafeCaches.buildNonEvictableCacheWithWeakInvalidateAll; +import static io.trino.execution.buffer.PagesSerdes.createSpillingPagesSerdeFactory; import static io.trino.spi.StandardErrorCode.OUT_OF_SPILL_SPACE; import static io.trino.util.Ciphers.createRandomAesEncryptionKey; import static java.lang.String.format; @@ -107,7 +108,7 @@ public FileSingleStreamSpillerFactory( CompressionCodec compressionCodec, boolean spillEncryptionEnabled) { - this.serdeFactory = new PagesSerdeFactory(blockEncodingSerde, compressionCodec); + this.serdeFactory = createSpillingPagesSerdeFactory(blockEncodingSerde, compressionCodec); this.executor = requireNonNull(executor, "executor is null"); this.spillerStats = requireNonNull(spillerStats, "spillerStats cannot be null"); requireNonNull(spillPaths, "spillPaths is null"); diff --git a/core/trino-main/src/main/java/io/trino/sql/planner/LocalExecutionPlanner.java b/core/trino-main/src/main/java/io/trino/sql/planner/LocalExecutionPlanner.java index 8c240a8a58fa..45f703adbfec 100644 --- a/core/trino-main/src/main/java/io/trino/sql/planner/LocalExecutionPlanner.java +++ b/core/trino-main/src/main/java/io/trino/sql/planner/LocalExecutionPlanner.java @@ -43,7 +43,6 @@ import io.trino.execution.TaskId; import io.trino.execution.TaskManagerConfig; import io.trino.execution.buffer.OutputBuffer; -import io.trino.execution.buffer.PagesSerdeFactory; import io.trino.metadata.MergeHandle; import io.trino.metadata.Metadata; import io.trino.metadata.ResolvedFunction; @@ -314,7 +313,6 @@ import static io.trino.SystemSessionProperties.getAdaptivePartialAggregationUniqueRowsRatioThreshold; import static io.trino.SystemSessionProperties.getAggregationOperatorUnspillMemoryLimit; import static io.trino.SystemSessionProperties.getDynamicRowFilterSelectivityThreshold; -import static io.trino.SystemSessionProperties.getExchangeCompressionCodec; import static io.trino.SystemSessionProperties.getFilterAndProjectMinOutputPageRowCount; import static io.trino.SystemSessionProperties.getFilterAndProjectMinOutputPageSize; import static io.trino.SystemSessionProperties.getPagePartitioningBufferPoolSize; @@ -331,6 +329,7 @@ import static io.trino.SystemSessionProperties.isSpillEnabled; import static io.trino.cache.CacheUtils.uncheckedCacheGet; import static io.trino.cache.SafeCaches.buildNonEvictableCache; +import static io.trino.execution.buffer.PagesSerdes.createExchangePagesSerdeFactory; import static io.trino.metadata.GlobalFunctionCatalog.builtinFunctionName; import static io.trino.operator.DistinctLimitOperator.DistinctLimitOperatorFactory; import static io.trino.operator.HashArraySizeSupplier.incrementalLoadFactorHashArraySizeSupplier; @@ -675,7 +674,7 @@ public LocalExecutionPlan plan( plan.getId(), outputTypes, pagePreprocessor, - new PagesSerdeFactory(plannerContext.getBlockEncodingSerde(), getExchangeCompressionCodec(session))), + createExchangePagesSerdeFactory(plannerContext.getBlockEncodingSerde(), session)), physicalOperation), context); @@ -940,7 +939,7 @@ private PhysicalOperation createMergeSource(RemoteSourceNode node, LocalExecutio context.getNextOperatorId(), node.getId(), directExchangeClientSupplier, - new PagesSerdeFactory(plannerContext.getBlockEncodingSerde(), getExchangeCompressionCodec(session)), + createExchangePagesSerdeFactory(plannerContext.getBlockEncodingSerde(), session), orderingCompiler, types, outputChannels, @@ -960,7 +959,7 @@ private PhysicalOperation createRemoteSource(RemoteSourceNode node, LocalExecuti context.getNextOperatorId(), node.getId(), directExchangeClientSupplier, - new PagesSerdeFactory(plannerContext.getBlockEncodingSerde(), getExchangeCompressionCodec(session)), + createExchangePagesSerdeFactory(plannerContext.getBlockEncodingSerde(), session), node.getRetryPolicy(), exchangeManagerRegistry); diff --git a/core/trino-main/src/test/java/io/trino/execution/TestSqlTaskExecution.java b/core/trino-main/src/test/java/io/trino/execution/TestSqlTaskExecution.java index 94b8c0c5dc67..ea8c722f42b4 100644 --- a/core/trino-main/src/test/java/io/trino/execution/TestSqlTaskExecution.java +++ b/core/trino-main/src/test/java/io/trino/execution/TestSqlTaskExecution.java @@ -31,7 +31,6 @@ import io.trino.execution.buffer.BufferState; import io.trino.execution.buffer.OutputBuffer; import io.trino.execution.buffer.OutputBufferStateMachine; -import io.trino.execution.buffer.PagesSerdeFactory; import io.trino.execution.buffer.PartitionedOutputBuffer; import io.trino.execution.buffer.PipelinedOutputBuffers; import io.trino.execution.buffer.PipelinedOutputBuffers.OutputBufferId; @@ -50,7 +49,6 @@ import io.trino.operator.output.TaskOutputOperator.TaskOutputOperatorFactory; import io.trino.spi.Page; import io.trino.spi.QueryId; -import io.trino.spi.block.TestingBlockEncodingSerde; import io.trino.spi.connector.ConnectorSplit; import io.trino.spiller.SpillSpaceTracker; import io.trino.sql.planner.LocalExecutionPlanner.LocalExecutionPlan; @@ -79,9 +77,9 @@ import static io.trino.execution.TaskState.RUNNING; import static io.trino.execution.TaskTestUtils.TABLE_SCAN_NODE_ID; import static io.trino.execution.TaskTestUtils.createTestSplitMonitor; -import static io.trino.execution.buffer.CompressionCodec.NONE; import static io.trino.execution.buffer.PagesSerdeUtil.getSerializedPagePositionCount; import static io.trino.execution.buffer.PipelinedOutputBuffers.BufferType.PARTITIONED; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.trino.testing.TestingHandles.TEST_CATALOG_HANDLE; import static java.util.Objects.requireNonNull; @@ -127,7 +125,7 @@ public void testSimple() TABLE_SCAN_NODE_ID, outputBuffer, Function.identity(), - new PagesSerdeFactory(new TestingBlockEncodingSerde(), NONE)); + createTestingPagesSerdeFactory()); LocalExecutionPlan localExecutionPlan = new LocalExecutionPlan( ImmutableList.of(new DriverFactory( 0, diff --git a/core/trino-main/src/test/java/io/trino/execution/buffer/BenchmarkBlockSerde.java b/core/trino-main/src/test/java/io/trino/execution/buffer/BenchmarkBlockSerde.java index 00eb240d9e4f..2d0dddb9adf5 100644 --- a/core/trino-main/src/test/java/io/trino/execution/buffer/BenchmarkBlockSerde.java +++ b/core/trino-main/src/test/java/io/trino/execution/buffer/BenchmarkBlockSerde.java @@ -23,7 +23,6 @@ import io.trino.spi.PageBuilder; import io.trino.spi.block.BlockBuilder; import io.trino.spi.block.RowBlockBuilder; -import io.trino.spi.block.TestingBlockEncodingSerde; import io.trino.spi.type.DecimalType; import io.trino.spi.type.Int128; import io.trino.spi.type.RowType; @@ -55,9 +54,9 @@ import static com.google.common.collect.ImmutableList.toImmutableList; import static io.airlift.slice.Slices.utf8Slice; import static io.trino.execution.buffer.BenchmarkDataGenerator.createValues; -import static io.trino.execution.buffer.CompressionCodec.NONE; import static io.trino.execution.buffer.PagesSerdeUtil.readPages; import static io.trino.execution.buffer.PagesSerdeUtil.writePages; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.jmh.Benchmarks.benchmark; import static io.trino.plugin.tpch.TpchTables.getTablePages; import static io.trino.spi.type.BigintType.BIGINT; @@ -207,7 +206,7 @@ public abstract static class TypeBenchmarkData public void setup(Type type, Function valueGenerator) { - PagesSerdeFactory serdeFactory = new PagesSerdeFactory(new TestingBlockEncodingSerde(), NONE); + PagesSerdeFactory serdeFactory = createTestingPagesSerdeFactory(); PageSerializer serializer = serdeFactory.createSerializer(Optional.empty()); PageDeserializer deserializer = serdeFactory.createDeserializer(Optional.empty()); PageBuilder pageBuilder = new PageBuilder(ImmutableList.of(type)); @@ -404,7 +403,7 @@ public static class LineitemBenchmarkData @Setup public void setup() { - PagesSerdeFactory serdeFactory = new PagesSerdeFactory(new TestingBlockEncodingSerde(), NONE); + PagesSerdeFactory serdeFactory = createTestingPagesSerdeFactory(); PageSerializer serializer = serdeFactory.createSerializer(Optional.empty()); PageDeserializer deserializer = serdeFactory.createDeserializer(Optional.empty()); diff --git a/core/trino-main/src/test/java/io/trino/execution/buffer/BenchmarkPagesSerde.java b/core/trino-main/src/test/java/io/trino/execution/buffer/BenchmarkPagesSerde.java index cd0044fcaace..88049640a2ca 100644 --- a/core/trino-main/src/test/java/io/trino/execution/buffer/BenchmarkPagesSerde.java +++ b/core/trino-main/src/test/java/io/trino/execution/buffer/BenchmarkPagesSerde.java @@ -18,7 +18,6 @@ import io.trino.spi.Page; import io.trino.spi.PageBuilder; import io.trino.spi.block.BlockBuilder; -import io.trino.spi.block.TestingBlockEncodingSerde; import io.trino.spi.type.Type; import org.junit.jupiter.api.Test; import org.openjdk.jmh.annotations.Benchmark; @@ -46,6 +45,7 @@ import static io.airlift.slice.Slices.utf8Slice; import static io.trino.execution.buffer.CompressionCodec.LZ4; import static io.trino.execution.buffer.CompressionCodec.NONE; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.jmh.Benchmarks.benchmark; import static io.trino.operator.PageAssertions.assertPageEquals; import static io.trino.spi.type.VarcharType.VARCHAR; @@ -115,7 +115,7 @@ public static class BenchmarkData @Setup public void initialize() { - PagesSerdeFactory serdeFactory = new PagesSerdeFactory(new TestingBlockEncodingSerde(), compressionCodec); + PagesSerdeFactory serdeFactory = createTestingPagesSerdeFactory(compressionCodec); Optional encryptionKey = encrypted ? Optional.of(createRandomAesEncryptionKey()) : Optional.empty(); serializer = serdeFactory.createSerializer(encryptionKey); deserializer = serdeFactory.createDeserializer(encryptionKey); diff --git a/core/trino-main/src/test/java/io/trino/execution/buffer/BufferTestUtils.java b/core/trino-main/src/test/java/io/trino/execution/buffer/BufferTestUtils.java index d990d046aaa4..353d27b01cc2 100644 --- a/core/trino-main/src/test/java/io/trino/execution/buffer/BufferTestUtils.java +++ b/core/trino-main/src/test/java/io/trino/execution/buffer/BufferTestUtils.java @@ -31,6 +31,8 @@ import static com.google.common.base.Preconditions.checkArgument; import static io.airlift.concurrent.MoreFutures.tryGetFutureValue; import static io.trino.execution.buffer.BufferState.FINISHED; +import static io.trino.execution.buffer.CompressionCodec.LZ4; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.SECONDS; import static org.assertj.core.api.Assertions.assertThat; @@ -39,7 +41,7 @@ public final class BufferTestUtils { private BufferTestUtils() {} - private static final PagesSerdeFactory PAGES_SERDE_FACTORY = new TestingPagesSerdeFactory(); + private static final PagesSerdeFactory PAGES_SERDE_FACTORY = createTestingPagesSerdeFactory(LZ4); static final Duration NO_WAIT = new Duration(0, MILLISECONDS); static final Duration MAX_WAIT = new Duration(1, SECONDS); private static final DataSize BUFFERED_PAGE_SIZE = DataSize.ofBytes(serializePage(createPage(42)).getRetainedSize()); diff --git a/core/trino-main/src/test/java/io/trino/execution/buffer/TestSpoolingExchangeOutputBuffer.java b/core/trino-main/src/test/java/io/trino/execution/buffer/TestSpoolingExchangeOutputBuffer.java index 8292ada11b3d..2e00ff8a0f4e 100644 --- a/core/trino-main/src/test/java/io/trino/execution/buffer/TestSpoolingExchangeOutputBuffer.java +++ b/core/trino-main/src/test/java/io/trino/execution/buffer/TestSpoolingExchangeOutputBuffer.java @@ -25,7 +25,6 @@ import io.trino.spi.Page; import io.trino.spi.PageBuilder; import io.trino.spi.QueryId; -import io.trino.spi.block.TestingBlockEncodingSerde; import io.trino.spi.block.VariableWidthBlockBuilder; import io.trino.spi.exchange.ExchangeSink; import io.trino.spi.exchange.ExchangeSinkInstanceHandle; @@ -42,7 +41,7 @@ import static io.trino.execution.buffer.BufferState.FINISHED; import static io.trino.execution.buffer.BufferState.FLUSHING; import static io.trino.execution.buffer.BufferState.NO_MORE_BUFFERS; -import static io.trino.execution.buffer.CompressionCodec.NONE; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.spi.type.VarcharType.VARCHAR; import static java.util.Objects.requireNonNull; import static org.assertj.core.api.Assertions.assertThat; @@ -311,7 +310,7 @@ private static Slice createPage(String value) VariableWidthBlockBuilder blockBuilder = (VariableWidthBlockBuilder) pageBuilder.getBlockBuilder(0); blockBuilder.writeEntry(valueSlice); Page page = pageBuilder.build(); - PageSerializer serializer = new PagesSerdeFactory(new TestingBlockEncodingSerde(), NONE).createSerializer(Optional.empty()); + PageSerializer serializer = createTestingPagesSerdeFactory().createSerializer(Optional.empty()); return serializer.serialize(page); } diff --git a/core/trino-main/src/test/java/io/trino/execution/buffer/TestingPagesSerdeFactory.java b/core/trino-main/src/test/java/io/trino/execution/buffer/TestingPagesSerdes.java similarity index 67% rename from core/trino-main/src/test/java/io/trino/execution/buffer/TestingPagesSerdeFactory.java rename to core/trino-main/src/test/java/io/trino/execution/buffer/TestingPagesSerdes.java index 3f3cea79b7f1..74b336922ace 100644 --- a/core/trino-main/src/test/java/io/trino/execution/buffer/TestingPagesSerdeFactory.java +++ b/core/trino-main/src/test/java/io/trino/execution/buffer/TestingPagesSerdes.java @@ -16,17 +16,22 @@ import io.trino.metadata.BlockEncodingManager; import io.trino.metadata.InternalBlockEncodingSerde; -import static io.trino.execution.buffer.CompressionCodec.LZ4; +import static io.trino.execution.buffer.CompressionCodec.NONE; import static io.trino.type.InternalTypeManager.TESTING_TYPE_MANAGER; -public class TestingPagesSerdeFactory - extends PagesSerdeFactory +public final class TestingPagesSerdes { + private TestingPagesSerdes() {} + private static final InternalBlockEncodingSerde BLOCK_ENCODING_SERDE = new InternalBlockEncodingSerde(new BlockEncodingManager(), TESTING_TYPE_MANAGER); - public TestingPagesSerdeFactory() + public static PagesSerdeFactory createTestingPagesSerdeFactory() + { + return createTestingPagesSerdeFactory(NONE); + } + + public static PagesSerdeFactory createTestingPagesSerdeFactory(CompressionCodec compressionCodec) { - // compression should be enabled in as many tests as possible - super(BLOCK_ENCODING_SERDE, LZ4); + return new PagesSerdeFactory(BLOCK_ENCODING_SERDE, compressionCodec); } } diff --git a/core/trino-main/src/test/java/io/trino/memory/TestMemoryPools.java b/core/trino-main/src/test/java/io/trino/memory/TestMemoryPools.java index 64efc87d0a2c..8bc6f85f7db4 100644 --- a/core/trino-main/src/test/java/io/trino/memory/TestMemoryPools.java +++ b/core/trino-main/src/test/java/io/trino/memory/TestMemoryPools.java @@ -20,7 +20,6 @@ import io.airlift.units.DataSize; import io.trino.execution.StageId; import io.trino.execution.TaskId; -import io.trino.execution.buffer.TestingPagesSerdeFactory; import io.trino.memory.context.LocalMemoryContext; import io.trino.operator.Driver; import io.trino.operator.DriverContext; @@ -52,6 +51,8 @@ import static io.airlift.units.DataSize.Unit.GIGABYTE; import static io.airlift.units.DataSize.Unit.MEGABYTE; import static io.trino.SessionTestUtils.TEST_SESSION; +import static io.trino.execution.buffer.CompressionCodec.LZ4; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.testing.TestingTaskContext.createTaskContext; import static java.lang.String.format; import static java.util.concurrent.Executors.newCachedThreadPool; @@ -100,7 +101,7 @@ private RevocableMemoryDriver createRevocableMemoryDriver(MemoryPool userPool, D TableScanOperator.class.getSimpleName()); OutputFactory outputFactory = new PageConsumerOutputFactory(types -> (page -> {})); - Operator outputOperator = outputFactory.createOutputOperator(2, new PlanNodeId("output"), ImmutableList.of(), Function.identity(), new TestingPagesSerdeFactory()).createOperator(driverContext); + Operator outputOperator = outputFactory.createOutputOperator(2, new PlanNodeId("output"), ImmutableList.of(), Function.identity(), createTestingPagesSerdeFactory(LZ4)).createOperator(driverContext); RevocableMemoryOperator revocableMemoryOperator = new RevocableMemoryOperator(revokableOperatorContext, reservedPerPage, numberOfPages); Driver driver = Driver.createDriver(driverContext, revocableMemoryOperator, outputOperator); diff --git a/core/trino-main/src/test/java/io/trino/operator/MockExchangeRequestProcessor.java b/core/trino-main/src/test/java/io/trino/operator/MockExchangeRequestProcessor.java index 46264bdd674b..e4a28bce944f 100644 --- a/core/trino-main/src/test/java/io/trino/operator/MockExchangeRequestProcessor.java +++ b/core/trino-main/src/test/java/io/trino/operator/MockExchangeRequestProcessor.java @@ -28,7 +28,6 @@ import io.trino.execution.buffer.BufferResult; import io.trino.execution.buffer.PageSerializer; import io.trino.execution.buffer.PagesSerdeFactory; -import io.trino.execution.buffer.TestingPagesSerdeFactory; import io.trino.server.InternalHeaders; import io.trino.spi.Page; @@ -47,7 +46,9 @@ import static com.google.common.net.HttpHeaders.CONTENT_TYPE; import static io.trino.TrinoMediaTypes.TRINO_PAGES; import static io.trino.cache.SafeCaches.buildNonEvictableCache; +import static io.trino.execution.buffer.CompressionCodec.LZ4; import static io.trino.execution.buffer.PagesSerdeUtil.calculateChecksum; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.server.InternalHeaders.TRINO_BUFFER_COMPLETE; import static io.trino.server.InternalHeaders.TRINO_PAGE_NEXT_TOKEN; import static io.trino.server.InternalHeaders.TRINO_PAGE_TOKEN; @@ -61,7 +62,7 @@ public class MockExchangeRequestProcessor { private static final String TASK_INSTANCE_ID = "task-instance-id"; - private final PagesSerdeFactory serdeFactory = new TestingPagesSerdeFactory(); + private final PagesSerdeFactory serdeFactory = createTestingPagesSerdeFactory(LZ4); private final LoadingCache buffers = buildNonEvictableCache(CacheBuilder.newBuilder(), CacheLoader.from(location -> new MockBuffer(location, serdeFactory.createSerializer(Optional.empty())))); diff --git a/core/trino-main/src/test/java/io/trino/operator/TestDirectExchangeClient.java b/core/trino-main/src/test/java/io/trino/operator/TestDirectExchangeClient.java index eb0a5a5475df..8ea56c3b50b8 100644 --- a/core/trino-main/src/test/java/io/trino/operator/TestDirectExchangeClient.java +++ b/core/trino-main/src/test/java/io/trino/operator/TestDirectExchangeClient.java @@ -39,7 +39,6 @@ import io.trino.execution.TaskId; import io.trino.execution.buffer.PageDeserializer; import io.trino.execution.buffer.PagesSerdeFactory; -import io.trino.execution.buffer.TestingPagesSerdeFactory; import io.trino.memory.context.SimpleLocalMemoryContext; import io.trino.spi.Page; import io.trino.spi.QueryId; @@ -76,7 +75,9 @@ import static io.airlift.concurrent.MoreFutures.tryGetFutureValue; import static io.airlift.concurrent.Threads.daemonThreadsNamed; import static io.trino.execution.TestSqlTaskExecution.TASK_ID; +import static io.trino.execution.buffer.CompressionCodec.LZ4; import static io.trino.execution.buffer.PagesSerdeUtil.getSerializedPagePositionCount; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.trino.spi.StandardErrorCode.GENERIC_INTERNAL_ERROR; import static io.trino.spi.exchange.ExchangeId.createRandomExchangeId; @@ -103,7 +104,7 @@ public void setUp() { scheduler = newScheduledThreadPool(4, daemonThreadsNamed(getClass().getSimpleName() + "-%s")); pageBufferClientCallbackExecutor = Executors.newSingleThreadExecutor(); - serdeFactory = new TestingPagesSerdeFactory(); + serdeFactory = createTestingPagesSerdeFactory(LZ4); } @AfterAll diff --git a/core/trino-main/src/test/java/io/trino/operator/TestExchangeOperator.java b/core/trino-main/src/test/java/io/trino/operator/TestExchangeOperator.java index 7fbe4e305536..f0b1f1d24796 100644 --- a/core/trino-main/src/test/java/io/trino/operator/TestExchangeOperator.java +++ b/core/trino-main/src/test/java/io/trino/operator/TestExchangeOperator.java @@ -31,7 +31,6 @@ import io.trino.execution.StageId; import io.trino.execution.TaskId; import io.trino.execution.buffer.PagesSerdeFactory; -import io.trino.execution.buffer.TestingPagesSerdeFactory; import io.trino.metadata.Split; import io.trino.operator.ExchangeOperator.ExchangeOperatorFactory; import io.trino.spi.Page; @@ -55,6 +54,8 @@ import static io.airlift.units.DataSize.Unit.MEGABYTE; import static io.trino.SessionTestUtils.TEST_SESSION; import static io.trino.cache.SafeCaches.buildNonEvictableCacheWithWeakInvalidateAll; +import static io.trino.execution.buffer.CompressionCodec.LZ4; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.operator.ExchangeOperator.REMOTE_CATALOG_HANDLE; import static io.trino.operator.PageAssertions.assertPageEquals; import static io.trino.operator.TestingTaskBuffer.PAGE; @@ -70,7 +71,7 @@ public class TestExchangeOperator { private static final List TYPES = ImmutableList.of(VARCHAR); - private static final PagesSerdeFactory SERDE_FACTORY = new TestingPagesSerdeFactory(); + private static final PagesSerdeFactory SERDE_FACTORY = createTestingPagesSerdeFactory(LZ4); private static final TaskId TASK_1_ID = new TaskId(new StageId("query", 0), 0, 0); private static final TaskId TASK_2_ID = new TaskId(new StageId("query", 0), 1, 0); diff --git a/core/trino-main/src/test/java/io/trino/operator/TestHttpPageBufferClient.java b/core/trino-main/src/test/java/io/trino/operator/TestHttpPageBufferClient.java index 46134f9836bc..e1dba8930d8c 100644 --- a/core/trino-main/src/test/java/io/trino/operator/TestHttpPageBufferClient.java +++ b/core/trino-main/src/test/java/io/trino/operator/TestHttpPageBufferClient.java @@ -29,7 +29,6 @@ import io.trino.execution.TaskId; import io.trino.execution.buffer.PageDeserializer; import io.trino.execution.buffer.PagesSerdeFactory; -import io.trino.execution.buffer.TestingPagesSerdeFactory; import io.trino.operator.HttpPageBufferClient.ClientCallback; import io.trino.spi.HostAddress; import io.trino.spi.Page; @@ -62,6 +61,8 @@ import static io.airlift.concurrent.Threads.daemonThreadsNamed; import static io.airlift.units.DataSize.Unit.MEGABYTE; import static io.trino.TrinoMediaTypes.TRINO_PAGES; +import static io.trino.execution.buffer.CompressionCodec.LZ4; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.spi.StandardErrorCode.EXCEEDED_LOCAL_MEMORY_LIMIT; import static io.trino.spi.StandardErrorCode.PAGE_TOO_LARGE; import static io.trino.spi.StandardErrorCode.PAGE_TRANSPORT_ERROR; @@ -559,7 +560,7 @@ private static void assertPageEquals(Page expectedPage, Page actualPage) private static class TestingClientCallback implements ClientCallback { - private final PagesSerdeFactory serdeFactory = new TestingPagesSerdeFactory(); + private final PagesSerdeFactory serdeFactory = createTestingPagesSerdeFactory(LZ4); private final CyclicBarrier done; private final List pages = Collections.synchronizedList(new ArrayList<>()); diff --git a/core/trino-main/src/test/java/io/trino/operator/TestMergeOperator.java b/core/trino-main/src/test/java/io/trino/operator/TestMergeOperator.java index a51ac83239d6..3523404c0f17 100644 --- a/core/trino-main/src/test/java/io/trino/operator/TestMergeOperator.java +++ b/core/trino-main/src/test/java/io/trino/operator/TestMergeOperator.java @@ -31,7 +31,6 @@ import io.trino.execution.StageId; import io.trino.execution.TaskId; import io.trino.execution.buffer.PagesSerdeFactory; -import io.trino.execution.buffer.TestingPagesSerdeFactory; import io.trino.metadata.Split; import io.trino.spi.Page; import io.trino.spi.connector.SortOrder; @@ -56,6 +55,8 @@ import static io.trino.RowPagesBuilder.rowPagesBuilder; import static io.trino.SessionTestUtils.TEST_SESSION; import static io.trino.cache.SafeCaches.buildNonEvictableCache; +import static io.trino.execution.buffer.CompressionCodec.LZ4; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.operator.OperatorAssertion.assertOperatorIsBlocked; import static io.trino.operator.OperatorAssertion.assertOperatorIsUnblocked; import static io.trino.operator.PageAssertions.assertPageEquals; @@ -89,7 +90,7 @@ public class TestMergeOperator public void setUp() { executor = newSingleThreadScheduledExecutor(daemonThreadsNamed("test-merge-operator-%s")); - serdeFactory = new TestingPagesSerdeFactory(); + serdeFactory = createTestingPagesSerdeFactory(LZ4); taskBuffers = buildNonEvictableCache(CacheBuilder.newBuilder(), CacheLoader.from(TestingTaskBuffer::new)); httpClient = new TestingHttpClient(new TestingExchangeHttpClientHandler(taskBuffers, serdeFactory), executor); diff --git a/core/trino-main/src/test/java/io/trino/operator/output/BenchmarkPartitionedOutputOperator.java b/core/trino-main/src/test/java/io/trino/operator/output/BenchmarkPartitionedOutputOperator.java index d445009e9217..70b58eb505d7 100644 --- a/core/trino-main/src/test/java/io/trino/operator/output/BenchmarkPartitionedOutputOperator.java +++ b/core/trino-main/src/test/java/io/trino/operator/output/BenchmarkPartitionedOutputOperator.java @@ -38,7 +38,6 @@ import io.trino.spi.block.Block; import io.trino.spi.block.RowBlock; import io.trino.spi.block.RunLengthEncodedBlock; -import io.trino.spi.block.TestingBlockEncodingSerde; import io.trino.spi.type.ArrayType; import io.trino.spi.type.BigintType; import io.trino.spi.type.BooleanType; @@ -93,6 +92,7 @@ import static io.trino.block.BlockAssertions.createRepeatedValuesBlock; import static io.trino.execution.buffer.CompressionCodec.NONE; import static io.trino.execution.buffer.PipelinedOutputBuffers.BufferType.PARTITIONED; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.trino.operator.output.BenchmarkPartitionedOutputOperator.BenchmarkData.TestType; import static io.trino.spi.type.BigintType.BIGINT; @@ -442,7 +442,7 @@ private PartitionedOutputOperator createPartitionedOutputOperator() PartitionFunction partitionFunction = new BucketPartitionFunction( new HashBucketFunction(new PrecomputedHashGenerator(0), partitionCount), IntStream.range(0, partitionCount).toArray()); - PagesSerdeFactory serdeFactory = new PagesSerdeFactory(new TestingBlockEncodingSerde(), compressionCodec); + PagesSerdeFactory serdeFactory = createTestingPagesSerdeFactory(compressionCodec); PartitionedOutputBuffer buffer = createPartitionedOutputBuffer(); diff --git a/core/trino-main/src/test/java/io/trino/operator/output/TestPagePartitioner.java b/core/trino-main/src/test/java/io/trino/operator/output/TestPagePartitioner.java index 53e502dff7fd..5acf08188db7 100644 --- a/core/trino-main/src/test/java/io/trino/operator/output/TestPagePartitioner.java +++ b/core/trino-main/src/test/java/io/trino/operator/output/TestPagePartitioner.java @@ -40,7 +40,6 @@ import io.trino.spi.block.Block; import io.trino.spi.block.DictionaryBlock; import io.trino.spi.block.RunLengthEncodedBlock; -import io.trino.spi.block.TestingBlockEncodingSerde; import io.trino.spi.predicate.NullableValue; import io.trino.spi.type.ArrayType; import io.trino.spi.type.Decimals; @@ -74,7 +73,7 @@ import static io.trino.block.BlockAssertions.createLongsBlock; import static io.trino.block.BlockAssertions.createRandomBlockForType; import static io.trino.block.BlockAssertions.createRepeatedValuesBlock; -import static io.trino.execution.buffer.CompressionCodec.NONE; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.trino.spi.type.BigintType.BIGINT; import static io.trino.spi.type.BooleanType.BOOLEAN; @@ -110,7 +109,7 @@ public class TestPagePartitioner private static final int POSITIONS_PER_PAGE = 8; private static final int PARTITION_COUNT = 2; - private static final PagesSerdeFactory PAGES_SERDE_FACTORY = new PagesSerdeFactory(new TestingBlockEncodingSerde(), NONE); + private static final PagesSerdeFactory PAGES_SERDE_FACTORY = createTestingPagesSerdeFactory(); private static final PageDeserializer PAGE_DESERIALIZER = PAGES_SERDE_FACTORY.createDeserializer(Optional.empty()); private final ExecutorService executor = newCachedThreadPool(daemonThreadsNamed(getClass().getSimpleName() + "-executor-%s")); diff --git a/core/trino-main/src/test/java/io/trino/operator/output/TestPagePartitionerPool.java b/core/trino-main/src/test/java/io/trino/operator/output/TestPagePartitionerPool.java index 97ccf2b444a4..276f1698e66b 100644 --- a/core/trino-main/src/test/java/io/trino/operator/output/TestPagePartitionerPool.java +++ b/core/trino-main/src/test/java/io/trino/operator/output/TestPagePartitionerPool.java @@ -25,7 +25,6 @@ import io.trino.execution.buffer.OutputBufferStatus; import io.trino.execution.buffer.OutputBuffers; import io.trino.execution.buffer.PipelinedOutputBuffers.OutputBufferId; -import io.trino.execution.buffer.TestingPagesSerdeFactory; import io.trino.memory.context.AggregatedMemoryContext; import io.trino.operator.BucketPartitionFunction; import io.trino.operator.DriverContext; @@ -55,6 +54,8 @@ import static io.airlift.concurrent.Threads.threadsNamed; import static io.trino.SessionTestUtils.TEST_SESSION; import static io.trino.block.BlockAssertions.createLongsBlock; +import static io.trino.execution.buffer.CompressionCodec.LZ4; +import static io.trino.execution.buffer.TestingPagesSerdes.createTestingPagesSerdeFactory; import static io.trino.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.trino.spi.type.BigintType.BIGINT; import static java.util.concurrent.Executors.newScheduledThreadPool; @@ -176,7 +177,7 @@ private static PartitionedOutputOperatorFactory createFactory(DataSize maxPagePa false, OptionalInt.empty(), outputBuffer, - new TestingPagesSerdeFactory(), + createTestingPagesSerdeFactory(LZ4), maxPagePartitioningBufferSize, new PositionsAppenderFactory(new BlockTypeOperators()), Optional.empty(), diff --git a/core/trino-main/src/test/java/io/trino/spiller/TestBinaryFileSpiller.java b/core/trino-main/src/test/java/io/trino/spiller/TestBinaryFileSpiller.java index 07bef0f650db..7eaf1e2ddc0e 100644 --- a/core/trino-main/src/test/java/io/trino/spiller/TestBinaryFileSpiller.java +++ b/core/trino-main/src/test/java/io/trino/spiller/TestBinaryFileSpiller.java @@ -42,6 +42,7 @@ import static com.google.common.io.MoreFiles.deleteRecursively; import static com.google.common.io.RecursiveDeleteOption.ALLOW_INSECURE; +import static io.trino.execution.buffer.PagesSerdes.createSpillingPagesSerdeFactory; import static io.trino.memory.context.AggregatedMemoryContext.newSimpleAggregatedMemoryContext; import static io.trino.operator.PageAssertions.assertPageEquals; import static io.trino.spi.type.BigintType.BIGINT; @@ -83,7 +84,7 @@ public void setUp() BlockEncodingSerde blockEncodingSerde = new TestingBlockEncodingSerde(); singleStreamSpillerFactory = new FileSingleStreamSpillerFactory(blockEncodingSerde, spillerStats, featuresConfig, nodeSpillConfig); factory = new GenericSpillerFactory(singleStreamSpillerFactory); - PagesSerdeFactory pagesSerdeFactory = new PagesSerdeFactory(blockEncodingSerde, nodeSpillConfig.getSpillCompressionCodec()); + PagesSerdeFactory pagesSerdeFactory = createSpillingPagesSerdeFactory(blockEncodingSerde, nodeSpillConfig.getSpillCompressionCodec()); serializer = pagesSerdeFactory.createSerializer(Optional.empty()); memoryContext = newSimpleAggregatedMemoryContext(); }