voonhous commented on code in PR #20112: URL: https://github.com/apache/hudi/pull/20112#discussion_r4203628518
########## 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: ACK, was gonna ask about CDC, make sense. -- 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]
