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]
