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]

Reply via email to