yihua commented on code in PR #20080:
URL: https://github.com/apache/hudi/pull/20080#discussion_r4188145977


##########
hudi-hadoop-common/src/test/java/org/apache/hudi/hadoop/fs/RecordingLocalFileSystem.java:
##########
@@ -0,0 +1,384 @@
+/*
+ * 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.hadoop.fs;
+
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.CreateFlag;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.LocatedFileStatus;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.fs.RawLocalFileSystem;
+import org.apache.hadoop.fs.RemoteIterator;
+import org.apache.hadoop.fs.permission.FsPermission;
+import org.apache.hadoop.util.Progressable;
+
+import java.io.IOException;
+import java.io.PrintWriter;
+import java.io.StringWriter;
+import java.util.Arrays;
+import java.util.EnumSet;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Queue;
+import java.util.WeakHashMap;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.function.BooleanSupplier;
+import java.util.function.Predicate;
+import java.util.stream.Collectors;
+
+/**
+ * The local file system with a recorder of every file system call, for tests 
that assert which
+ * calls a code path makes: how many times it lists the timeline, whether it 
opens a file once per
+ * task or once per job, whether an engine task touches the {@code .hoodie} 
folder at all.
+ *
+ * <p>Each call is recorded with its operation, path and thread, and whether 
it ran in the current
+ * scope, a test-defined notion such as "inside a Spark task" or "inside the 
compaction operator"
+ * (see {@link #setScope}). Calls in scope also keep their call stack. Only 
the outermost call is
+ * recorded, so {@code exists} delegating to {@code getFileStatus} counts 
once. Compound operations
+ * that a store serves as several requests ({@code globStatus}, {@code 
listFiles},
+ * {@code listStatusIterator}) are not recorded themselves; the directory 
listings and status calls
+ * they make are, when they make them. Tests filter the recorded calls with 
the predicates on
+ * {@link Call}, for example {@code 
getCalls(Call.inScope().and(Call.underMetaFolder()))}.
+ *
+ * <p>Usage: {@link #register} the file system on the Hadoop configuration the 
code under test uses
+ * (with Flink options, set {@code "hadoop." + FILE_IMPL_KEY} and {@code 
"hadoop." + DISABLE_CACHE_KEY}
+ * instead), {@link #reset}, run the work and inspect {@link #getCalls}. 
Assert that a call the work
+ * must make was recorded before asserting that another one was not, so that 
the test fails rather
+ * than passes when the recorder is not in use.
+ *
+ * <p>The recorder state is static because Hadoop creates a new instance per 
{@code FileSystem.get}
+ * call once caching is disabled, which is needed for the implementation to be 
picked up regardless
+ * of what the file system cache already holds.
+ */
+public class RecordingLocalFileSystem extends RawLocalFileSystem {

Review Comment:
   Done, the class doc now says every test in the JVM shares the recorder 
state, and that it extends `RawLocalFileSystem` (it writes no `.crc` files, and 
listings include ones the stock file system wrote).



##########
hudi-hadoop-common/src/test/java/org/apache/hudi/hadoop/fs/RecordingLocalFileSystem.java:
##########
@@ -0,0 +1,384 @@
+/*
+ * 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.hadoop.fs;
+
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.CreateFlag;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.LocatedFileStatus;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.fs.RawLocalFileSystem;
+import org.apache.hadoop.fs.RemoteIterator;
+import org.apache.hadoop.fs.permission.FsPermission;
+import org.apache.hadoop.util.Progressable;
+
+import java.io.IOException;
+import java.io.PrintWriter;
+import java.io.StringWriter;
+import java.util.Arrays;
+import java.util.EnumSet;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Queue;
+import java.util.WeakHashMap;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.function.BooleanSupplier;
+import java.util.function.Predicate;
+import java.util.stream.Collectors;
+
+/**
+ * The local file system with a recorder of every file system call, for tests 
that assert which
+ * calls a code path makes: how many times it lists the timeline, whether it 
opens a file once per
+ * task or once per job, whether an engine task touches the {@code .hoodie} 
folder at all.
+ *
+ * <p>Each call is recorded with its operation, path and thread, and whether 
it ran in the current
+ * scope, a test-defined notion such as "inside a Spark task" or "inside the 
compaction operator"
+ * (see {@link #setScope}). Calls in scope also keep their call stack. Only 
the outermost call is
+ * recorded, so {@code exists} delegating to {@code getFileStatus} counts 
once. Compound operations
+ * that a store serves as several requests ({@code globStatus}, {@code 
listFiles},
+ * {@code listStatusIterator}) are not recorded themselves; the directory 
listings and status calls
+ * they make are, when they make them. Tests filter the recorded calls with 
the predicates on
+ * {@link Call}, for example {@code 
getCalls(Call.inScope().and(Call.underMetaFolder()))}.
+ *
+ * <p>Usage: {@link #register} the file system on the Hadoop configuration the 
code under test uses
+ * (with Flink options, set {@code "hadoop." + FILE_IMPL_KEY} and {@code 
"hadoop." + DISABLE_CACHE_KEY}
+ * instead), {@link #reset}, run the work and inspect {@link #getCalls}. 
Assert that a call the work
+ * must make was recorded before asserting that another one was not, so that 
the test fails rather
+ * than passes when the recorder is not in use.
+ *
+ * <p>The recorder state is static because Hadoop creates a new instance per 
{@code FileSystem.get}
+ * call once caching is disabled, which is needed for the implementation to be 
picked up regardless
+ * of what the file system cache already holds.
+ */
+public class RecordingLocalFileSystem extends RawLocalFileSystem {
+
+  public static final String FILE_IMPL_KEY = "fs.file.impl";
+  public static final String DISABLE_CACHE_KEY = "fs.file.impl.disable.cache";
+
+  private static final String META_FOLDER = "/.hoodie";
+  private static final int MAX_DESCRIBED_STACKS = 3;
+  private static final BooleanSupplier NO_SCOPE = () -> false;
+
+  private static final Queue<Call> CALLS = new ConcurrentLinkedQueue<>();
+  private static final ThreadLocal<Boolean> IN_RECORDED_CALL = 
ThreadLocal.withInitial(() -> false);
+  private static final Map<Configuration, Map<String, String>> 
VALUES_BEFORE_REGISTER = new WeakHashMap<>();
+  private static volatile BooleanSupplier scope = NO_SCOPE;
+
+  /**
+   * One file system call.
+   */
+  public static final class Call {
+    private final String operation;
+    private final String path;
+    private final String threadName;
+    private final boolean inScope;
+    private final boolean underMetaFolder;
+    private final Throwable stack;
+
+    private Call(String operation, Path path, boolean inScope) {
+      this.operation = operation;
+      this.path = String.valueOf(path);
+      this.threadName = Thread.currentThread().getName();
+      this.inScope = inScope;
+      this.underMetaFolder = isMetaFolderPath(path);
+      this.stack = inScope ? new Throwable(operation + " " + path) : null;
+    }
+
+    public String getOperation() {
+      return operation;
+    }
+
+    public String getPath() {
+      return path;
+    }
+
+    public String getThreadName() {
+      return threadName;
+    }
+
+    public boolean isInScope() {
+      return inScope;
+    }
+
+    public boolean isUnderMetaFolder() {
+      return underMetaFolder;
+    }
+
+    /**
+     * The call stack, kept for calls in scope only.
+     */
+    public Throwable getStack() {
+      return stack;
+    }
+
+    public static Predicate<Call> inScope() {
+      return Call::isInScope;
+    }
+
+    /**
+     * Calls to a table's {@code .hoodie} folder or anything under it, the 
metadata table included.
+     */
+    public static Predicate<Call> underMetaFolder() {
+      return Call::isUnderMetaFolder;
+    }
+
+    public static Predicate<Call> operation(String... operations) {
+      List<String> names = Arrays.asList(operations);
+      return call -> names.contains(call.getOperation());
+    }
+
+    public static Predicate<Call> pathEndsWith(String suffix) {
+      return call -> call.getPath().endsWith(suffix);
+    }
+
+    public static Predicate<Call> pathContains(String part) {
+      return call -> call.getPath().contains(part);
+    }
+
+    @Override
+    public String toString() {
+      return operation + " " + path + " [" + threadName + "]";
+    }
+  }
+
+  /**
+   * Restores the scope that was set before {@link #withScope}.
+   */
+  public static final class Scope implements AutoCloseable {
+    private final BooleanSupplier previous;
+
+    private Scope(BooleanSupplier previous) {
+      this.previous = previous;
+    }
+
+    @Override
+    public void close() {
+      scope = previous;
+    }
+  }
+
+  /**
+   * Makes the {@code file} scheme resolve to this file system on the given 
configuration.
+   */
+  public static void register(Configuration conf) {
+    synchronized (VALUES_BEFORE_REGISTER) {
+      VALUES_BEFORE_REGISTER.computeIfAbsent(conf, c -> {
+        Map<String, String> values = new HashMap<>();
+        values.put(FILE_IMPL_KEY, c.getRaw(FILE_IMPL_KEY));
+        values.put(DISABLE_CACHE_KEY, c.getRaw(DISABLE_CACHE_KEY));
+        return values;
+      });
+    }
+    conf.setClass(FILE_IMPL_KEY, RecordingLocalFileSystem.class, 
FileSystem.class);
+    conf.setBoolean(DISABLE_CACHE_KEY, true);

Review Comment:
   The class doc already explains this in its paragraph on why the recorder 
state is static, so I'd rather not repeat it at the call site.



##########
hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestFileSystemBackedTableMetadata.java:
##########
@@ -229,6 +241,40 @@ public void testStrayFilesAreFilteredOut() throws 
Exception {
         "Stray files should be filtered out by listPartitions");
   }
 
+  /**
+   * The meta folder is never a partition, so listing must neither probe it 
for partition metadata
+   * nor list it as a partition, even when it holds a stray partition metafile.
+   */
+  @Test
+  void testMetaFolderIsNotProbedForPartitionMetadata() throws Exception {
+    hoodieTestTable = hoodieTestTable.addCommit("100");
+    for (String partition : ONE_LEVEL_PARTITIONS) {
+      hoodieTestTable = hoodieTestTable.withPartitionMetaFiles(partition)
+          .withBaseFilesInPartition(partition, IntStream.range(0, 
2).toArray());
+    }
+    try (OutputStream out = metaClient.getStorage().create(
+        new StoragePath(metaClient.getMetaPath(), 
HoodiePartitionMetadata.HOODIE_PARTITION_METAFILE_PREFIX))) {
+      out.write("stray".getBytes());
+    }
+
+    StorageConfiguration<Configuration> conf = 
HadoopFSUtils.getStorageConfWithCopy(
+        metaClient.getStorageConf().unwrapAs(Configuration.class));
+    RecordingLocalFileSystem.register(conf.unwrap());
+    RecordingLocalFileSystem.reset();

Review Comment:
   `register` only changes a throwaway copy of the configuration, and the 
helper's usage contract has each test call `reset()` before it runs, so an 
`@AfterEach` would change nothing here.



-- 
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