yihua commented on code in PR #20112: URL: https://github.com/apache/hudi/pull/20112#discussion_r4201715814
########## hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestSparkReadExecutorFootprint.java: ########## @@ -0,0 +1,434 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hudi.functional; + +import org.apache.hudi.DataSourceReadOptions; +import org.apache.hudi.DataSourceWriteOptions; +import org.apache.hudi.SparkAdapterSupport$; +import org.apache.hudi.common.config.HoodieMetadataConfig; +import org.apache.hudi.common.config.HoodieStorageConfig; +import org.apache.hudi.common.model.HoodieTableType; +import org.apache.hudi.common.table.HoodieTableConfig; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.HoodieTableVersion; +import org.apache.hudi.common.table.timeline.HoodieInstant; +import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.testutils.SparkClientFunctionalTestHarness; +import org.apache.hudi.testutils.SparkExecutorGuards; +import org.apache.hudi.testutils.TaskDeserializationRecorder; + +import lombok.extern.slf4j.Slf4j; +import org.apache.spark.SparkConf; +import org.apache.spark.sql.DataFrameReader; +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.RowFactory; +import org.apache.spark.sql.SaveMode; +import org.apache.spark.sql.types.DataTypes; +import org.apache.spark.sql.types.StructField; +import org.apache.spark.sql.types.StructType; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import java.util.stream.Stream; + +import static org.apache.hudi.common.model.HoodieTableType.COPY_ON_WRITE; +import static org.apache.hudi.common.model.HoodieTableType.MERGE_ON_READ; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Guards what Spark tasks do on the executors when reading a Hudi table through the file group + * reader: they must not access the table's {@code .hoodie} folder, they must not deserialize heavy + * driver-side objects (meta client, timeline, Hadoop configuration, the file format itself) with + * every task, what every task deserializes must stay within a size budget, and they must not parse + * the Hadoop default resources for every file they read. These costs scale with the number of tasks + * or files, not with the data. + * + * <p>Each table is written once per class and read by every case that needs it. The table has + * several partitions so that a read runs several tasks; on MERGE_ON_READ the second commit + * updates half of the keys so that every file group has log files to merge. + */ +@Slf4j +@Tag("functional") +class TestSparkReadExecutorFootprint extends SparkClientFunctionalTestHarness { + + private static final int CURRENT_VERSION = HoodieTableVersion.current().versionCode(); + private static final int NUM_RECORDS = 200; + private static final int NUM_UPDATED_RECORDS = 100; + private static final int NUM_PARTITIONS = 4; + + + /** + * Driver-side classes that a read task must not deserialize with its closure. + */ + private static final List<String> CLASSES_NOT_DESERIALIZED_PER_TASK = Arrays.asList( + "org.apache.hudi.common.table.HoodieTableMetaClient", + "org.apache.hudi.common.table.timeline.HoodieActiveTimeline", + "org.apache.hudi.storage.StorageConfiguration", + "org.apache.spark.util.SerializableConfiguration", + "org.apache.spark.sql.execution.datasources.parquet.HoodieFileGroupReaderBasedFileFormat", + "org.apache.hudi.config.HoodieWriteConfig"); + + /** + * Hadoop default resource parses allowed in the tasks of a read. A base file read converts the table + * schema in the scan state with a new configuration once per JVM instance of the state, so once per + * executor, and local mode runs one executor; a merging read does not parse them. + */ + private static final int MAX_HADOOP_DEFAULT_RESOURCE_LOADS = 1; + + private static final StructType SCHEMA = DataTypes.createStructType(new StructField[] { + DataTypes.createStructField("key", DataTypes.StringType, false), + DataTypes.createStructField("part", DataTypes.StringType, false), + DataTypes.createStructField("ts", DataTypes.LongType, false), + DataTypes.createStructField("value", DataTypes.StringType, true)}); + + private static final String VECTORIZED_READER_ENABLED = "spark.sql.parquet.enableVectorizedReader"; + + private static final Map<String, TestTable> TABLES = new HashMap<>(); + + @TempDir + static Path tablesDir; + + enum TableKind { + PLAIN, CDC, PARQUET_LOG_BLOCKS + } + + /** + * What every task of a read may deserialize: the task binary, the Java-serialized closure, and the + * largest task stream, which is the larger of the binary and the task with its partition. Any state + * added to the closure or to a task's partition is paid again by every task, so the budgets are about + * 1.4x the sizes measured on Spark 3.5, whatever the state is. A base file read carries the columnar + * reader in its closure, so it gets a larger budget than a merging read. + */ + enum TaskBudget { + // measured: task binary 17555 bytes, largest task stream 18647 bytes + BASE_FILE_READ(25 * 1024, 26 * 1024), + // measured: task binary 10205 to 11051 bytes, largest task stream 12225 to 12488 bytes + MERGING_READ(16 * 1024, 18 * 1024); + + private final long maxTaskBinaryBytes; + private final long maxTaskStreamBytes; + + TaskBudget(long maxTaskBinaryBytes, long maxTaskStreamBytes) { + this.maxTaskBinaryBytes = maxTaskBinaryBytes; + this.maxTaskStreamBytes = maxTaskStreamBytes; + } + } Review Comment: Done. The budgets are now kept per Spark and Scala binary version (3.3, 3.4, 3.5 on Scala 2.12 and 2.13, 4.0, 4.1, 4.2), each about 1.1x the largest size measured on that version, separately for columnar and row reads. A Spark version without a budget fails the test until one is added. ########## hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestSparkReadExecutorFootprint.java: ########## @@ -0,0 +1,434 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hudi.functional; + +import org.apache.hudi.DataSourceReadOptions; +import org.apache.hudi.DataSourceWriteOptions; +import org.apache.hudi.SparkAdapterSupport$; +import org.apache.hudi.common.config.HoodieMetadataConfig; +import org.apache.hudi.common.config.HoodieStorageConfig; +import org.apache.hudi.common.model.HoodieTableType; +import org.apache.hudi.common.table.HoodieTableConfig; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.HoodieTableVersion; +import org.apache.hudi.common.table.timeline.HoodieInstant; +import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.testutils.SparkClientFunctionalTestHarness; +import org.apache.hudi.testutils.SparkExecutorGuards; +import org.apache.hudi.testutils.TaskDeserializationRecorder; + +import lombok.extern.slf4j.Slf4j; +import org.apache.spark.SparkConf; +import org.apache.spark.sql.DataFrameReader; +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.RowFactory; +import org.apache.spark.sql.SaveMode; +import org.apache.spark.sql.types.DataTypes; +import org.apache.spark.sql.types.StructField; +import org.apache.spark.sql.types.StructType; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import java.util.stream.Stream; + +import static org.apache.hudi.common.model.HoodieTableType.COPY_ON_WRITE; +import static org.apache.hudi.common.model.HoodieTableType.MERGE_ON_READ; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Guards what Spark tasks do on the executors when reading a Hudi table through the file group + * reader: they must not access the table's {@code .hoodie} folder, they must not deserialize heavy + * driver-side objects (meta client, timeline, Hadoop configuration, the file format itself) with + * every task, what every task deserializes must stay within a size budget, and they must not parse + * the Hadoop default resources for every file they read. These costs scale with the number of tasks + * or files, not with the data. + * + * <p>Each table is written once per class and read by every case that needs it. The table has + * several partitions so that a read runs several tasks; on MERGE_ON_READ the second commit + * updates half of the keys so that every file group has log files to merge. + */ +@Slf4j +@Tag("functional") +class TestSparkReadExecutorFootprint extends SparkClientFunctionalTestHarness { + + private static final int CURRENT_VERSION = HoodieTableVersion.current().versionCode(); + private static final int NUM_RECORDS = 200; + private static final int NUM_UPDATED_RECORDS = 100; + private static final int NUM_PARTITIONS = 4; + + + /** + * Driver-side classes that a read task must not deserialize with its closure. + */ + private static final List<String> CLASSES_NOT_DESERIALIZED_PER_TASK = Arrays.asList( + "org.apache.hudi.common.table.HoodieTableMetaClient", + "org.apache.hudi.common.table.timeline.HoodieActiveTimeline", + "org.apache.hudi.storage.StorageConfiguration", + "org.apache.spark.util.SerializableConfiguration", + "org.apache.spark.sql.execution.datasources.parquet.HoodieFileGroupReaderBasedFileFormat", + "org.apache.hudi.config.HoodieWriteConfig"); + + /** + * Hadoop default resource parses allowed in the tasks of a read. A base file read converts the table + * schema in the scan state with a new configuration once per JVM instance of the state, so once per + * executor, and local mode runs one executor; a merging read does not parse them. + */ + private static final int MAX_HADOOP_DEFAULT_RESOURCE_LOADS = 1; + + private static final StructType SCHEMA = DataTypes.createStructType(new StructField[] { + DataTypes.createStructField("key", DataTypes.StringType, false), + DataTypes.createStructField("part", DataTypes.StringType, false), + DataTypes.createStructField("ts", DataTypes.LongType, false), + DataTypes.createStructField("value", DataTypes.StringType, true)}); + + private static final String VECTORIZED_READER_ENABLED = "spark.sql.parquet.enableVectorizedReader"; + + private static final Map<String, TestTable> TABLES = new HashMap<>(); + + @TempDir + static Path tablesDir; + + enum TableKind { + PLAIN, CDC, PARQUET_LOG_BLOCKS + } + + /** + * What every task of a read may deserialize: the task binary, the Java-serialized closure, and the + * largest task stream, which is the larger of the binary and the task with its partition. Any state + * added to the closure or to a task's partition is paid again by every task, so the budgets are about + * 1.4x the sizes measured on Spark 3.5, whatever the state is. A base file read carries the columnar + * reader in its closure, so it gets a larger budget than a merging read. + */ + enum TaskBudget { + // measured: task binary 17555 bytes, largest task stream 18647 bytes + BASE_FILE_READ(25 * 1024, 26 * 1024), + // measured: task binary 10205 to 11051 bytes, largest task stream 12225 to 12488 bytes + MERGING_READ(16 * 1024, 18 * 1024); + + private final long maxTaskBinaryBytes; + private final long maxTaskStreamBytes; + + TaskBudget(long maxTaskBinaryBytes, long maxTaskStreamBytes) { + this.maxTaskBinaryBytes = maxTaskBinaryBytes; + this.maxTaskStreamBytes = maxTaskStreamBytes; + } + } + + enum ReadQuery { + SNAPSHOT, READ_OPTIMIZED, INCREMENTAL, TIME_TRAVEL, CDC + } + + @Override + public SparkConf conf() { + return conf(Collections.singletonMap("spark.plugins", SparkExecutorGuards.TASK_START_HOOK_PLUGIN)); + } + + @BeforeEach + void enableRecording() { + // Until #20090, a read can leave the session's vectorized reader flag changed; reset it so that + // every case plans its scan the same way whatever ran before it. + spark().conf().unset(VECTORIZED_READER_ENABLED); + SparkExecutorGuards.enableFileSystemCallRecording(jsc().hadoopConfiguration()); + } + + @AfterEach + void disableRecording() { + SparkExecutorGuards.disableFileSystemCallRecording(jsc().hadoopConfiguration()); + } + + @AfterAll + static void forgetTables() { + TABLES.clear(); + } + + static Stream<Arguments> tableVersionsAndTypes() { + return Stream.of(6, CURRENT_VERSION).flatMap(version -> + Stream.of(COPY_ON_WRITE, MERGE_ON_READ).map(type -> Arguments.of(version, type))); + } + + /** + * With the metadata table off, a query other than an incremental one lists partitions from the + * file system with a Spark job, whose tasks are guarded too. + */ + static Stream<Arguments> readQueries() { + List<Arguments> args = new ArrayList<>(); + for (int version : new int[] {6, CURRENT_VERSION}) { + for (HoodieTableType type : HoodieTableType.values()) { + for (ReadQuery query : Arrays.asList(ReadQuery.SNAPSHOT, ReadQuery.READ_OPTIMIZED, ReadQuery.INCREMENTAL, ReadQuery.TIME_TRAVEL)) { + if (query == ReadQuery.READ_OPTIMIZED && type == COPY_ON_WRITE) { + continue; + } + for (String metadataOnRead : new String[] {"default", "false"}) { + args.add(Arguments.of(version, type, query, metadataOnRead)); + } + } + } + } + return args.stream(); + } + + /** + * A CDC read of a version 6 MERGE_ON_READ table fails on the driver, so it is left out. + */ + static Stream<Arguments> cdcTableVersionsAndTypes() { + return tableVersionsAndTypes().filter(args -> !isVersion6MergeOnRead(args)); + } + + private static boolean isVersion6MergeOnRead(Arguments args) { + Object[] values = args.get(); + return (int) values[0] == 6 && values[1] == MERGE_ON_READ; + } + + static Stream<Arguments> readsForHadoopDefaultResources() { + return Stream.of(6, CURRENT_VERSION).flatMap(version -> Stream.of( + Arguments.of(version, COPY_ON_WRITE, TableKind.PLAIN), + Arguments.of(version, MERGE_ON_READ, TableKind.PLAIN), + Arguments.of(version, MERGE_ON_READ, TableKind.PARQUET_LOG_BLOCKS))); + } + + static Stream<Arguments> readsForDeserialization() { + return Stream.of( + Arguments.of(CURRENT_VERSION, COPY_ON_WRITE, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.BASE_FILE_READ), + Arguments.of(6, COPY_ON_WRITE, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.BASE_FILE_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.MERGING_READ), + Arguments.of(6, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.READ_OPTIMIZED, TaskBudget.BASE_FILE_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.INCREMENTAL, TaskBudget.MERGING_READ), + Arguments.of(6, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.INCREMENTAL, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.TIME_TRAVEL, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.CDC, ReadQuery.CDC, TaskBudget.MERGING_READ)); + } + + @ParameterizedTest(name = "[{index}] version={0}, type={1}, query={2}, metadata={3}") + @MethodSource("readQueries") + void testNoExecutorMetaFolderAccess(int tableVersion, HoodieTableType tableType, ReadQuery query, String metadataOnRead) { + TestTable table = getOrWriteTable(tableVersion, tableType, TableKind.PLAIN); + Map<String, String> options = new HashMap<>(); + if (!"default".equals(metadataOnRead)) { + options.put(HoodieMetadataConfig.ENABLE.key(), metadataOnRead); + } + List<Row> rows = SparkExecutorGuards.assertNoExecutorMetaFolderAccess( + table.name + " " + query + " read with metadata " + metadataOnRead, + () -> read(table, query, options).collectAsList()); + assertFalse(rows.isEmpty(), "The read must return rows for the guard to be meaningful"); + } + + @ParameterizedTest(name = "[{index}] version={0}, type={1}") + @MethodSource("cdcTableVersionsAndTypes") + void testNoExecutorMetaFolderAccessForCdcQuery(int tableVersion, HoodieTableType tableType) { + TestTable table = getOrWriteTable(tableVersion, tableType, TableKind.CDC); + List<Row> rows = SparkExecutorGuards.assertNoExecutorMetaFolderAccess( + table.name + " CDC read", + () -> read(table, ReadQuery.CDC, new HashMap<>()).collectAsList()); + assertFalse(rows.isEmpty(), "The CDC read must return rows for the guard to be meaningful"); + } + + /** + * Tasks must not deserialize the meta client, timeline, Hadoop configuration, write config or the + * file format with their closure, and what every task deserializes must stay within budget. + */ + @ParameterizedTest(name = "[{index}] version={0}, type={1}, table={2}, query={3}, budget={4}") + @MethodSource("readsForDeserialization") + void testTaskDeserializationFootprint(int tableVersion, HoodieTableType tableType, TableKind kind, ReadQuery query, + TaskBudget budget) { + TestTable table = getOrWriteTable(tableVersion, tableType, kind); + Dataset<Row> df = read(table, query, new HashMap<>()); + // Plan and list files on the driver first, so that the recorded window holds only the scan. + df.queryExecution().executedPlan().execute(); + List<Row> rows = new ArrayList<>(); + TaskDeserializationRecorder.Result result = SparkExecutorGuards.recordTaskDeserialization( + spark().sparkContext(), () -> rows.addAll(df.collectAsList())); + assertFalse(rows.isEmpty(), "The read must return rows for the guard to be meaningful"); + SparkExecutorGuards.TaskBinary taskBinary = SparkExecutorGuards.inspectTaskBinary(df); + log.info("{} read of {}: task binary {} bytes, largest task stream seen {} bytes, stages kept {}, ignored {}", + query, table.name, taskBinary.getBytes(), result.getMaxStreamBytes(), result.getKeptScopes(), result.getIgnoredScopes()); + SparkExecutorGuards.assertTaskDeserializationFootprint( + table.name + " " + query + " read", result, taskBinary, CLASSES_NOT_DESERIALIZED_PER_TASK, + budget.maxTaskBinaryBytes, budget.maxTaskStreamBytes); + } + + /** + * Tasks must read base files and log blocks with the configuration they are given rather than + * create one per file, which parses the Hadoop default resources every time. + */ + @ParameterizedTest(name = "[{index}] version={0}, type={1}, table={2}") + @MethodSource("readsForHadoopDefaultResources") + void testNoHadoopDefaultResourceLoadsPerFile(int tableVersion, HoodieTableType tableType, TableKind kind) { + TestTable table = getOrWriteTable(tableVersion, tableType, kind); + Dataset<Row> df = read(table, ReadQuery.SNAPSHOT, new HashMap<>()); + List<Row> rows = SparkExecutorGuards.assertTaskHadoopDefaultResourceLoadsAtMost( + table.name + " snapshot read", spark().sparkContext(), MAX_HADOOP_DEFAULT_RESOURCE_LOADS, df::collectAsList); + assertEquals(NUM_RECORDS, rows.size()); + } + + /** + * With data skipping, file pruning consults the column stats in the metadata table. The tasks of + * the scan must still not read the timeline or table config of either table. + */ + @ParameterizedTest(name = "[{index}] version={0}, type={1}") + @MethodSource("tableVersionsAndTypes") + void testNoExecutorMetaFolderAccessWithDataSkipping(int tableVersion, HoodieTableType tableType) { + TestTable table = getOrWriteTable(tableVersion, tableType, TableKind.PLAIN); + Map<String, String> options = new HashMap<>(); + options.put(DataSourceReadOptions.ENABLE_DATA_SKIPPING().key(), "true"); + List<Row> rows = SparkExecutorGuards.assertNoExecutorMetaFolderAccess( + table.name + " snapshot read with data skipping", + () -> read(table, ReadQuery.SNAPSHOT, options).filter("value = 'v2'").collectAsList()); + assertEquals(NUM_UPDATED_RECORDS, rows.size()); + } Review Comment: Done. The read now looks up the first key in strict mode. Every partition holds the keys with the same `i % 4`, so column stats prune the files of three of the four partitions, and the test asserts the tasks read only files of the first one. ########## hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestSparkReadExecutorFootprint.java: ########## @@ -0,0 +1,434 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hudi.functional; + +import org.apache.hudi.DataSourceReadOptions; +import org.apache.hudi.DataSourceWriteOptions; +import org.apache.hudi.SparkAdapterSupport$; +import org.apache.hudi.common.config.HoodieMetadataConfig; +import org.apache.hudi.common.config.HoodieStorageConfig; +import org.apache.hudi.common.model.HoodieTableType; +import org.apache.hudi.common.table.HoodieTableConfig; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.HoodieTableVersion; +import org.apache.hudi.common.table.timeline.HoodieInstant; +import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.testutils.SparkClientFunctionalTestHarness; +import org.apache.hudi.testutils.SparkExecutorGuards; +import org.apache.hudi.testutils.TaskDeserializationRecorder; + +import lombok.extern.slf4j.Slf4j; +import org.apache.spark.SparkConf; +import org.apache.spark.sql.DataFrameReader; +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.RowFactory; +import org.apache.spark.sql.SaveMode; +import org.apache.spark.sql.types.DataTypes; +import org.apache.spark.sql.types.StructField; +import org.apache.spark.sql.types.StructType; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import java.util.stream.Stream; + +import static org.apache.hudi.common.model.HoodieTableType.COPY_ON_WRITE; +import static org.apache.hudi.common.model.HoodieTableType.MERGE_ON_READ; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Guards what Spark tasks do on the executors when reading a Hudi table through the file group + * reader: they must not access the table's {@code .hoodie} folder, they must not deserialize heavy + * driver-side objects (meta client, timeline, Hadoop configuration, the file format itself) with + * every task, what every task deserializes must stay within a size budget, and they must not parse + * the Hadoop default resources for every file they read. These costs scale with the number of tasks + * or files, not with the data. + * + * <p>Each table is written once per class and read by every case that needs it. The table has + * several partitions so that a read runs several tasks; on MERGE_ON_READ the second commit + * updates half of the keys so that every file group has log files to merge. + */ +@Slf4j +@Tag("functional") +class TestSparkReadExecutorFootprint extends SparkClientFunctionalTestHarness { + + private static final int CURRENT_VERSION = HoodieTableVersion.current().versionCode(); + private static final int NUM_RECORDS = 200; + private static final int NUM_UPDATED_RECORDS = 100; + private static final int NUM_PARTITIONS = 4; + + + /** + * Driver-side classes that a read task must not deserialize with its closure. + */ + private static final List<String> CLASSES_NOT_DESERIALIZED_PER_TASK = Arrays.asList( + "org.apache.hudi.common.table.HoodieTableMetaClient", + "org.apache.hudi.common.table.timeline.HoodieActiveTimeline", + "org.apache.hudi.storage.StorageConfiguration", + "org.apache.spark.util.SerializableConfiguration", + "org.apache.spark.sql.execution.datasources.parquet.HoodieFileGroupReaderBasedFileFormat", + "org.apache.hudi.config.HoodieWriteConfig"); + + /** + * Hadoop default resource parses allowed in the tasks of a read. A base file read converts the table + * schema in the scan state with a new configuration once per JVM instance of the state, so once per + * executor, and local mode runs one executor; a merging read does not parse them. + */ + private static final int MAX_HADOOP_DEFAULT_RESOURCE_LOADS = 1; + + private static final StructType SCHEMA = DataTypes.createStructType(new StructField[] { + DataTypes.createStructField("key", DataTypes.StringType, false), + DataTypes.createStructField("part", DataTypes.StringType, false), + DataTypes.createStructField("ts", DataTypes.LongType, false), + DataTypes.createStructField("value", DataTypes.StringType, true)}); + + private static final String VECTORIZED_READER_ENABLED = "spark.sql.parquet.enableVectorizedReader"; + + private static final Map<String, TestTable> TABLES = new HashMap<>(); + + @TempDir + static Path tablesDir; + + enum TableKind { + PLAIN, CDC, PARQUET_LOG_BLOCKS + } + + /** + * What every task of a read may deserialize: the task binary, the Java-serialized closure, and the + * largest task stream, which is the larger of the binary and the task with its partition. Any state + * added to the closure or to a task's partition is paid again by every task, so the budgets are about + * 1.4x the sizes measured on Spark 3.5, whatever the state is. A base file read carries the columnar + * reader in its closure, so it gets a larger budget than a merging read. + */ + enum TaskBudget { + // measured: task binary 17555 bytes, largest task stream 18647 bytes + BASE_FILE_READ(25 * 1024, 26 * 1024), + // measured: task binary 10205 to 11051 bytes, largest task stream 12225 to 12488 bytes + MERGING_READ(16 * 1024, 18 * 1024); + + private final long maxTaskBinaryBytes; + private final long maxTaskStreamBytes; + + TaskBudget(long maxTaskBinaryBytes, long maxTaskStreamBytes) { + this.maxTaskBinaryBytes = maxTaskBinaryBytes; + this.maxTaskStreamBytes = maxTaskStreamBytes; + } + } + + enum ReadQuery { + SNAPSHOT, READ_OPTIMIZED, INCREMENTAL, TIME_TRAVEL, CDC + } + + @Override + public SparkConf conf() { + return conf(Collections.singletonMap("spark.plugins", SparkExecutorGuards.TASK_START_HOOK_PLUGIN)); + } + + @BeforeEach + void enableRecording() { + // Until #20090, a read can leave the session's vectorized reader flag changed; reset it so that + // every case plans its scan the same way whatever ran before it. + spark().conf().unset(VECTORIZED_READER_ENABLED); + SparkExecutorGuards.enableFileSystemCallRecording(jsc().hadoopConfiguration()); + } + + @AfterEach + void disableRecording() { + SparkExecutorGuards.disableFileSystemCallRecording(jsc().hadoopConfiguration()); + } + + @AfterAll + static void forgetTables() { + TABLES.clear(); + } + + static Stream<Arguments> tableVersionsAndTypes() { + return Stream.of(6, CURRENT_VERSION).flatMap(version -> + Stream.of(COPY_ON_WRITE, MERGE_ON_READ).map(type -> Arguments.of(version, type))); + } + + /** + * With the metadata table off, a query other than an incremental one lists partitions from the + * file system with a Spark job, whose tasks are guarded too. + */ + static Stream<Arguments> readQueries() { + List<Arguments> args = new ArrayList<>(); + for (int version : new int[] {6, CURRENT_VERSION}) { + for (HoodieTableType type : HoodieTableType.values()) { + for (ReadQuery query : Arrays.asList(ReadQuery.SNAPSHOT, ReadQuery.READ_OPTIMIZED, ReadQuery.INCREMENTAL, ReadQuery.TIME_TRAVEL)) { + if (query == ReadQuery.READ_OPTIMIZED && type == COPY_ON_WRITE) { + continue; + } + for (String metadataOnRead : new String[] {"default", "false"}) { + args.add(Arguments.of(version, type, query, metadataOnRead)); + } + } + } + } + return args.stream(); + } + + /** + * A CDC read of a version 6 MERGE_ON_READ table fails on the driver, so it is left out. + */ + static Stream<Arguments> cdcTableVersionsAndTypes() { + return tableVersionsAndTypes().filter(args -> !isVersion6MergeOnRead(args)); + } + + private static boolean isVersion6MergeOnRead(Arguments args) { + Object[] values = args.get(); + return (int) values[0] == 6 && values[1] == MERGE_ON_READ; + } + + static Stream<Arguments> readsForHadoopDefaultResources() { + return Stream.of(6, CURRENT_VERSION).flatMap(version -> Stream.of( + Arguments.of(version, COPY_ON_WRITE, TableKind.PLAIN), + Arguments.of(version, MERGE_ON_READ, TableKind.PLAIN), + Arguments.of(version, MERGE_ON_READ, TableKind.PARQUET_LOG_BLOCKS))); + } + + static Stream<Arguments> readsForDeserialization() { + return Stream.of( + Arguments.of(CURRENT_VERSION, COPY_ON_WRITE, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.BASE_FILE_READ), + Arguments.of(6, COPY_ON_WRITE, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.BASE_FILE_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.MERGING_READ), + Arguments.of(6, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.READ_OPTIMIZED, TaskBudget.BASE_FILE_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.INCREMENTAL, TaskBudget.MERGING_READ), + Arguments.of(6, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.INCREMENTAL, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.TIME_TRAVEL, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.CDC, ReadQuery.CDC, TaskBudget.MERGING_READ)); + } + + @ParameterizedTest(name = "[{index}] version={0}, type={1}, query={2}, metadata={3}") + @MethodSource("readQueries") + void testNoExecutorMetaFolderAccess(int tableVersion, HoodieTableType tableType, ReadQuery query, String metadataOnRead) { + TestTable table = getOrWriteTable(tableVersion, tableType, TableKind.PLAIN); + Map<String, String> options = new HashMap<>(); + if (!"default".equals(metadataOnRead)) { + options.put(HoodieMetadataConfig.ENABLE.key(), metadataOnRead); + } + List<Row> rows = SparkExecutorGuards.assertNoExecutorMetaFolderAccess( + table.name + " " + query + " read with metadata " + metadataOnRead, + () -> read(table, query, options).collectAsList()); + assertFalse(rows.isEmpty(), "The read must return rows for the guard to be meaningful"); + } + + @ParameterizedTest(name = "[{index}] version={0}, type={1}") + @MethodSource("cdcTableVersionsAndTypes") + void testNoExecutorMetaFolderAccessForCdcQuery(int tableVersion, HoodieTableType tableType) { + TestTable table = getOrWriteTable(tableVersion, tableType, TableKind.CDC); + List<Row> rows = SparkExecutorGuards.assertNoExecutorMetaFolderAccess( + table.name + " CDC read", + () -> read(table, ReadQuery.CDC, new HashMap<>()).collectAsList()); + assertFalse(rows.isEmpty(), "The CDC read must return rows for the guard to be meaningful"); + } + + /** + * Tasks must not deserialize the meta client, timeline, Hadoop configuration, write config or the + * file format with their closure, and what every task deserializes must stay within budget. + */ + @ParameterizedTest(name = "[{index}] version={0}, type={1}, table={2}, query={3}, budget={4}") + @MethodSource("readsForDeserialization") + void testTaskDeserializationFootprint(int tableVersion, HoodieTableType tableType, TableKind kind, ReadQuery query, + TaskBudget budget) { + TestTable table = getOrWriteTable(tableVersion, tableType, kind); + Dataset<Row> df = read(table, query, new HashMap<>()); + // Plan and list files on the driver first, so that the recorded window holds only the scan. + df.queryExecution().executedPlan().execute(); + List<Row> rows = new ArrayList<>(); + TaskDeserializationRecorder.Result result = SparkExecutorGuards.recordTaskDeserialization( + spark().sparkContext(), () -> rows.addAll(df.collectAsList())); + assertFalse(rows.isEmpty(), "The read must return rows for the guard to be meaningful"); + SparkExecutorGuards.TaskBinary taskBinary = SparkExecutorGuards.inspectTaskBinary(df); + log.info("{} read of {}: task binary {} bytes, largest task stream seen {} bytes, stages kept {}, ignored {}", + query, table.name, taskBinary.getBytes(), result.getMaxStreamBytes(), result.getKeptScopes(), result.getIgnoredScopes()); + SparkExecutorGuards.assertTaskDeserializationFootprint( + table.name + " " + query + " read", result, taskBinary, CLASSES_NOT_DESERIALIZED_PER_TASK, + budget.maxTaskBinaryBytes, budget.maxTaskStreamBytes); + } + + /** + * Tasks must read base files and log blocks with the configuration they are given rather than + * create one per file, which parses the Hadoop default resources every time. + */ + @ParameterizedTest(name = "[{index}] version={0}, type={1}, table={2}") + @MethodSource("readsForHadoopDefaultResources") + void testNoHadoopDefaultResourceLoadsPerFile(int tableVersion, HoodieTableType tableType, TableKind kind) { + TestTable table = getOrWriteTable(tableVersion, tableType, kind); + Dataset<Row> df = read(table, ReadQuery.SNAPSHOT, new HashMap<>()); + List<Row> rows = SparkExecutorGuards.assertTaskHadoopDefaultResourceLoadsAtMost( + table.name + " snapshot read", spark().sparkContext(), MAX_HADOOP_DEFAULT_RESOURCE_LOADS, df::collectAsList); + assertEquals(NUM_RECORDS, rows.size()); + } + + /** + * With data skipping, file pruning consults the column stats in the metadata table. The tasks of + * the scan must still not read the timeline or table config of either table. + */ + @ParameterizedTest(name = "[{index}] version={0}, type={1}") + @MethodSource("tableVersionsAndTypes") + void testNoExecutorMetaFolderAccessWithDataSkipping(int tableVersion, HoodieTableType tableType) { + TestTable table = getOrWriteTable(tableVersion, tableType, TableKind.PLAIN); + Map<String, String> options = new HashMap<>(); + options.put(DataSourceReadOptions.ENABLE_DATA_SKIPPING().key(), "true"); + List<Row> rows = SparkExecutorGuards.assertNoExecutorMetaFolderAccess( + table.name + " snapshot read with data skipping", + () -> read(table, ReadQuery.SNAPSHOT, options).filter("value = 'v2'").collectAsList()); + assertEquals(NUM_UPDATED_RECORDS, rows.size()); + } + + private Dataset<Row> read(TestTable table, ReadQuery query, Map<String, String> options) { + DataFrameReader reader = spark().read().format("hudi").options(options); + switch (query) { + case SNAPSHOT: + reader.option(DataSourceReadOptions.QUERY_TYPE().key(), DataSourceReadOptions.QUERY_TYPE_SNAPSHOT_OPT_VAL()); + break; + case READ_OPTIMIZED: + reader.option(DataSourceReadOptions.QUERY_TYPE().key(), DataSourceReadOptions.QUERY_TYPE_READ_OPTIMIZED_OPT_VAL()); + break; + case TIME_TRAVEL: + reader.option(DataSourceReadOptions.TIME_TRAVEL_AS_OF_INSTANT().key(), table.lastInstant.requestedTime()); Review Comment: Done. Time travel now reads as of the first commit, and the test asserts it sees only the first commit's values. ########## hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestSparkReadExecutorFootprint.java: ########## @@ -0,0 +1,434 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hudi.functional; + +import org.apache.hudi.DataSourceReadOptions; +import org.apache.hudi.DataSourceWriteOptions; +import org.apache.hudi.SparkAdapterSupport$; +import org.apache.hudi.common.config.HoodieMetadataConfig; +import org.apache.hudi.common.config.HoodieStorageConfig; +import org.apache.hudi.common.model.HoodieTableType; +import org.apache.hudi.common.table.HoodieTableConfig; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.HoodieTableVersion; +import org.apache.hudi.common.table.timeline.HoodieInstant; +import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.testutils.SparkClientFunctionalTestHarness; +import org.apache.hudi.testutils.SparkExecutorGuards; +import org.apache.hudi.testutils.TaskDeserializationRecorder; + +import lombok.extern.slf4j.Slf4j; +import org.apache.spark.SparkConf; +import org.apache.spark.sql.DataFrameReader; +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.RowFactory; +import org.apache.spark.sql.SaveMode; +import org.apache.spark.sql.types.DataTypes; +import org.apache.spark.sql.types.StructField; +import org.apache.spark.sql.types.StructType; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import java.util.stream.Stream; + +import static org.apache.hudi.common.model.HoodieTableType.COPY_ON_WRITE; +import static org.apache.hudi.common.model.HoodieTableType.MERGE_ON_READ; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Guards what Spark tasks do on the executors when reading a Hudi table through the file group + * reader: they must not access the table's {@code .hoodie} folder, they must not deserialize heavy + * driver-side objects (meta client, timeline, Hadoop configuration, the file format itself) with + * every task, what every task deserializes must stay within a size budget, and they must not parse + * the Hadoop default resources for every file they read. These costs scale with the number of tasks + * or files, not with the data. + * + * <p>Each table is written once per class and read by every case that needs it. The table has + * several partitions so that a read runs several tasks; on MERGE_ON_READ the second commit + * updates half of the keys so that every file group has log files to merge. + */ +@Slf4j +@Tag("functional") +class TestSparkReadExecutorFootprint extends SparkClientFunctionalTestHarness { + + private static final int CURRENT_VERSION = HoodieTableVersion.current().versionCode(); + private static final int NUM_RECORDS = 200; + private static final int NUM_UPDATED_RECORDS = 100; + private static final int NUM_PARTITIONS = 4; + + + /** + * Driver-side classes that a read task must not deserialize with its closure. + */ + private static final List<String> CLASSES_NOT_DESERIALIZED_PER_TASK = Arrays.asList( + "org.apache.hudi.common.table.HoodieTableMetaClient", + "org.apache.hudi.common.table.timeline.HoodieActiveTimeline", + "org.apache.hudi.storage.StorageConfiguration", + "org.apache.spark.util.SerializableConfiguration", + "org.apache.spark.sql.execution.datasources.parquet.HoodieFileGroupReaderBasedFileFormat", + "org.apache.hudi.config.HoodieWriteConfig"); + + /** + * Hadoop default resource parses allowed in the tasks of a read. A base file read converts the table + * schema in the scan state with a new configuration once per JVM instance of the state, so once per + * executor, and local mode runs one executor; a merging read does not parse them. + */ + private static final int MAX_HADOOP_DEFAULT_RESOURCE_LOADS = 1; + + private static final StructType SCHEMA = DataTypes.createStructType(new StructField[] { + DataTypes.createStructField("key", DataTypes.StringType, false), + DataTypes.createStructField("part", DataTypes.StringType, false), + DataTypes.createStructField("ts", DataTypes.LongType, false), + DataTypes.createStructField("value", DataTypes.StringType, true)}); + + private static final String VECTORIZED_READER_ENABLED = "spark.sql.parquet.enableVectorizedReader"; + + private static final Map<String, TestTable> TABLES = new HashMap<>(); + + @TempDir + static Path tablesDir; + + enum TableKind { + PLAIN, CDC, PARQUET_LOG_BLOCKS + } + + /** + * What every task of a read may deserialize: the task binary, the Java-serialized closure, and the + * largest task stream, which is the larger of the binary and the task with its partition. Any state + * added to the closure or to a task's partition is paid again by every task, so the budgets are about + * 1.4x the sizes measured on Spark 3.5, whatever the state is. A base file read carries the columnar + * reader in its closure, so it gets a larger budget than a merging read. + */ + enum TaskBudget { + // measured: task binary 17555 bytes, largest task stream 18647 bytes + BASE_FILE_READ(25 * 1024, 26 * 1024), + // measured: task binary 10205 to 11051 bytes, largest task stream 12225 to 12488 bytes + MERGING_READ(16 * 1024, 18 * 1024); + + private final long maxTaskBinaryBytes; + private final long maxTaskStreamBytes; + + TaskBudget(long maxTaskBinaryBytes, long maxTaskStreamBytes) { + this.maxTaskBinaryBytes = maxTaskBinaryBytes; + this.maxTaskStreamBytes = maxTaskStreamBytes; + } + } + + enum ReadQuery { + SNAPSHOT, READ_OPTIMIZED, INCREMENTAL, TIME_TRAVEL, CDC + } + + @Override + public SparkConf conf() { + return conf(Collections.singletonMap("spark.plugins", SparkExecutorGuards.TASK_START_HOOK_PLUGIN)); + } + + @BeforeEach + void enableRecording() { + // Until #20090, a read can leave the session's vectorized reader flag changed; reset it so that + // every case plans its scan the same way whatever ran before it. + spark().conf().unset(VECTORIZED_READER_ENABLED); + SparkExecutorGuards.enableFileSystemCallRecording(jsc().hadoopConfiguration()); + } + + @AfterEach + void disableRecording() { + SparkExecutorGuards.disableFileSystemCallRecording(jsc().hadoopConfiguration()); + } + + @AfterAll + static void forgetTables() { + TABLES.clear(); + } + + static Stream<Arguments> tableVersionsAndTypes() { + return Stream.of(6, CURRENT_VERSION).flatMap(version -> + Stream.of(COPY_ON_WRITE, MERGE_ON_READ).map(type -> Arguments.of(version, type))); + } + + /** + * With the metadata table off, a query other than an incremental one lists partitions from the + * file system with a Spark job, whose tasks are guarded too. + */ + static Stream<Arguments> readQueries() { + List<Arguments> args = new ArrayList<>(); + for (int version : new int[] {6, CURRENT_VERSION}) { + for (HoodieTableType type : HoodieTableType.values()) { + for (ReadQuery query : Arrays.asList(ReadQuery.SNAPSHOT, ReadQuery.READ_OPTIMIZED, ReadQuery.INCREMENTAL, ReadQuery.TIME_TRAVEL)) { + if (query == ReadQuery.READ_OPTIMIZED && type == COPY_ON_WRITE) { + continue; + } + for (String metadataOnRead : new String[] {"default", "false"}) { + args.add(Arguments.of(version, type, query, metadataOnRead)); + } + } + } + } + return args.stream(); + } + + /** + * A CDC read of a version 6 MERGE_ON_READ table fails on the driver, so it is left out. + */ + static Stream<Arguments> cdcTableVersionsAndTypes() { + return tableVersionsAndTypes().filter(args -> !isVersion6MergeOnRead(args)); + } + + private static boolean isVersion6MergeOnRead(Arguments args) { + Object[] values = args.get(); + return (int) values[0] == 6 && values[1] == MERGE_ON_READ; + } + + static Stream<Arguments> readsForHadoopDefaultResources() { + return Stream.of(6, CURRENT_VERSION).flatMap(version -> Stream.of( + Arguments.of(version, COPY_ON_WRITE, TableKind.PLAIN), + Arguments.of(version, MERGE_ON_READ, TableKind.PLAIN), + Arguments.of(version, MERGE_ON_READ, TableKind.PARQUET_LOG_BLOCKS))); + } + + static Stream<Arguments> readsForDeserialization() { + return Stream.of( + Arguments.of(CURRENT_VERSION, COPY_ON_WRITE, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.BASE_FILE_READ), + Arguments.of(6, COPY_ON_WRITE, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.BASE_FILE_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.MERGING_READ), + Arguments.of(6, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.SNAPSHOT, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.READ_OPTIMIZED, TaskBudget.BASE_FILE_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.INCREMENTAL, TaskBudget.MERGING_READ), + Arguments.of(6, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.INCREMENTAL, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.PLAIN, ReadQuery.TIME_TRAVEL, TaskBudget.MERGING_READ), + Arguments.of(CURRENT_VERSION, MERGE_ON_READ, TableKind.CDC, ReadQuery.CDC, TaskBudget.MERGING_READ)); + } Review Comment: Done. The size rows are now generated for every query on both table versions and types, except the version 6 MERGE_ON_READ CDC read, which fails on the driver until #20071; its guard is in #20208. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
