This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 75d6dfb49f [core] Parallelize per-partition file listing for 
catalog-managed format tables (#8845)
75d6dfb49f is described below

commit 75d6dfb49fecbc94c381d7983103ce8bb212a226
Author: XiaoHongbo <[email protected]>
AuthorDate: Sun Jul 26 15:12:10 2026 +0800

    [core] Parallelize per-partition file listing for catalog-managed format 
tables (#8845)
---
 docs/generated/core_configuration.html             |   6 +
 .../main/java/org/apache/paimon/CoreOptions.java   |  12 +
 .../table/format/CatalogSplitEnumerator.java       | 252 ++++++++++++++++
 .../table/format/FileSystemSplitEnumerator.java    | 132 +++++++++
 .../paimon/table/format/FormatTableScan.java       | 327 +--------------------
 .../paimon/table/format/SplitEnumerator.java       | 168 +++++++++++
 .../format/CatalogManagedPartitionScanTest.java    | 282 ++++++++++++++++++
 7 files changed, 859 insertions(+), 320 deletions(-)

diff --git a/docs/generated/core_configuration.html 
b/docs/generated/core_configuration.html
index 2841618e7c..d6f44e9f33 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -746,6 +746,12 @@ under the License.
             <td>Boolean</td>
             <td>Format table file path only contain partition value.</td>
         </tr>
+        <tr>
+            <td><h5>format-table.scan.list-parallelism</h5></td>
+            <td style="word-wrap: break-word;">64</td>
+            <td>Integer</td>
+            <td>The parallelism of listing partition files during split 
planning for a Format Table with catalog-managed partitions.</td>
+        </tr>
         <tr>
             <td><h5>full-compaction.delta-commits</h5></td>
             <td style="word-wrap: break-word;">(none)</td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java 
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index e0c7fd4004..c5989ff9b3 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -1230,6 +1230,14 @@ public class CoreOptions implements Serializable {
                                     + "Note: Scale-up this parameter will 
increase memory usage while scanning manifest files. "
                                     + "We can consider downsize it when we 
encounter an out of memory exception while scanning");
 
+    public static final ConfigOption<Integer> 
FORMAT_TABLE_SCAN_LIST_PARALLELISM =
+            key("format-table.scan.list-parallelism")
+                    .intType()
+                    .defaultValue(64)
+                    .withDescription(
+                            "The parallelism of listing partition files during 
split planning for "
+                                    + "a Format Table with catalog-managed 
partitions.");
+
     public static final ConfigOption<Duration> STREAMING_READ_SNAPSHOT_DELAY =
             key("streaming.read.snapshot.delay")
                     .durationType()
@@ -3708,6 +3716,10 @@ public class CoreOptions implements Serializable {
         return options.get(SCAN_MANIFEST_PARALLELISM);
     }
 
+    public Integer formatTableScanListParallelism() {
+        return options.get(FORMAT_TABLE_SCAN_LIST_PARALLELISM);
+    }
+
     public Integer scanBucket() {
         return options.get(SCAN_BUCKET);
     }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/format/CatalogSplitEnumerator.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/format/CatalogSplitEnumerator.java
new file mode 100644
index 0000000000..eaa43e3d69
--- /dev/null
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/format/CatalogSplitEnumerator.java
@@ -0,0 +1,252 @@
+/*
+ * 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.paimon.table.format;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.fs.FileIO;
+import org.apache.paimon.fs.FileStatus;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.manifest.PartitionEntry;
+import org.apache.paimon.partition.Partition;
+import org.apache.paimon.partition.PartitionPredicate;
+import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.table.FormatTable;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.utils.Pair;
+import org.apache.paimon.utils.PartitionPathUtils;
+import org.apache.paimon.utils.SemaphoredDelegatingExecutor;
+import org.apache.paimon.utils.ThreadPoolUtils;
+import org.apache.paimon.utils.ThreadUtils;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import javax.annotation.Nullable;
+
+import java.io.FileNotFoundException;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.SynchronousQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Function;
+
+/** A {@link SplitEnumerator} whose partitions are managed by the catalog. */
+final class CatalogSplitEnumerator extends SplitEnumerator {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(CatalogSplitEnumerator.class);
+
+    private static final int LIST_POOL_MAX_THREADS = 1000;
+
+    // Cached pool bounded at 1000 threads that reuses idle workers; 
CallerRunsPolicy lists on the
+    // caller for back pressure once the cap is hit.
+    private static final ThreadPoolExecutor LIST_POOL =
+            new ThreadPoolExecutor(
+                    0,
+                    LIST_POOL_MAX_THREADS,
+                    1,
+                    TimeUnit.MINUTES,
+                    new SynchronousQueue<>(),
+                    
ThreadUtils.newDaemonThreadFactory("FORMAT-TABLE-LIST-THREAD-POOL"),
+                    new ThreadPoolExecutor.CallerRunsPolicy());
+
+    private final FormatTablePartitionManager partitionManager;
+
+    CatalogSplitEnumerator(
+            FormatTable table,
+            CoreOptions coreOptions,
+            FormatTablePartitionManager partitionManager) {
+        super(table, coreOptions);
+        this.partitionManager = partitionManager;
+    }
+
+    @Override
+    List<Split> enumeratePartitions(@Nullable PartitionPredicate 
partitionFilter)
+            throws IOException {
+        List<Pair<LinkedHashMap<String, String>, Path>> partitions =
+                findPartitions(partitionFilter);
+        List<Split> splits = new ArrayList<>();
+        if (partitions.isEmpty()) {
+            return splits;
+        }
+
+        FileIO fileIO = table.fileIO();
+        // Establish the filesystem on the caller thread so listing workers 
reuse it under the
+        // caller's security context instead of creating it lazily under a 
shared worker.
+        fileIO.exists(new Path(table.location()));
+        Function<Pair<LinkedHashMap<String, String>, Path>, List<Split>> 
lister =
+                pair -> {
+                    BinaryRow partitionRow = toPartitionRow(pair.getKey());
+                    if (partitionFilter != null && 
!partitionFilter.test(partitionRow)) {
+                        return Collections.emptyList();
+                    }
+                    try {
+                        return createSplits(fileIO, pair.getValue(), 
partitionRow);
+                    } catch (FileNotFoundException e) {
+                        warnMissingPartition(pair.getKey(), pair.getValue());
+                        return Collections.emptyList();
+                    } catch (IOException e) {
+                        throw new RuntimeException(
+                                "Failed to list files for partition " + 
pair.getValue(), e);
+                    }
+                };
+        int parallelism =
+                Math.min(
+                        LIST_POOL_MAX_THREADS,
+                        Math.max(1, 
coreOptions.formatTableScanListParallelism()));
+        ExecutorService executor = new SemaphoredDelegatingExecutor(LIST_POOL, 
parallelism, false);
+        ThreadPoolUtils.randomlyExecuteSequentialReturn(executor, lister, 
partitions)
+                .forEachRemaining(splits::add);
+        return splits;
+    }
+
+    @Override
+    List<Pair<LinkedHashMap<String, String>, Path>> findPartitions(
+            @Nullable PartitionPredicate partitionFilter) {
+        Optional<Predicate> extracted = 
FormatTableScan.extractPartitionPredicate(partitionFilter);
+        Map<String, String> prefix = leadingEqualityPrefix(extracted);
+        Predicate catalogFilter = extracted.orElse(null);
+        List<Partition> partitions = partitionManager.listPartitions(prefix, 
catalogFilter);
+        if (partitions.isEmpty() && prefix.isEmpty() && catalogFilter == null) 
{
+            warnIfFilesystemPartitionsExist();
+        }
+        return toSpecsAndPaths(partitions, 
coreOptions.formatTablePartitionOnlyValueInPath());
+    }
+
+    @Override
+    List<PartitionEntry> listPartitionEntries() {
+        List<Partition> partitions = 
partitionManager.listPartitions(Collections.emptyMap(), null);
+        if (partitions.isEmpty()) {
+            warnIfFilesystemPartitionsExist();
+        }
+        boolean onlyValueInPath = 
coreOptions.formatTablePartitionOnlyValueInPath();
+        List<PartitionEntry> entries = new ArrayList<>(partitions.size());
+        Set<Map<String, String>> seen = new HashSet<>(partitions.size());
+        for (Partition partition : partitions) {
+            if (!seen.add(partition.spec())) {
+                continue;
+            }
+            entries.add(
+                    new PartitionEntry(
+                            toPartitionRow(normalizeSpec(partition.spec(), 
onlyValueInPath)),
+                            partition.recordCount(),
+                            partition.fileSizeInBytes(),
+                            partition.fileCount(),
+                            partition.lastFileCreationTime(),
+                            partition.totalBuckets()));
+        }
+        return entries;
+    }
+
+    private List<Pair<LinkedHashMap<String, String>, Path>> toSpecsAndPaths(
+            List<Partition> partitions, boolean onlyValueInPath) {
+        List<Pair<LinkedHashMap<String, String>, Path>> result = new 
ArrayList<>(partitions.size());
+        Path tablePath = new Path(table.location());
+        // A duplicate catalog entry must not duplicate all records in that 
partition.
+        Set<String> seenPartitionPaths = new HashSet<>(partitions.size());
+        for (Partition partition : partitions) {
+            LinkedHashMap<String, String> spec = 
normalizeSpec(partition.spec(), onlyValueInPath);
+            String partitionPath =
+                    PartitionPathUtils.generatePartitionPathUtil(spec, 
onlyValueInPath);
+            if (seenPartitionPaths.add(partitionPath)) {
+                result.add(Pair.of(spec, new Path(tablePath, partitionPath)));
+            }
+        }
+        return result;
+    }
+
+    LinkedHashMap<String, String> normalizeSpec(
+            @Nullable Map<String, String> spec, boolean onlyValueInPath) {
+        List<String> partitionKeys = table.partitionKeys();
+        if (spec == null
+                || spec.size() != partitionKeys.size()
+                || !spec.keySet().containsAll(partitionKeys)) {
+            throw corruptPartitionSpec(spec);
+        }
+        LinkedHashMap<String, String> normalized = new LinkedHashMap<>();
+        for (String partitionKey : partitionKeys) {
+            String value = spec.get(partitionKey);
+            // In a value-only layout, "." and ".." would resolve outside the 
table.
+            try {
+                PartitionPathUtils.validatePartitionValueForPath(value, 
onlyValueInPath);
+            } catch (IllegalArgumentException e) {
+                throw corruptPartitionSpec(spec);
+            }
+            normalized.put(partitionKey, value);
+        }
+        return normalized;
+    }
+
+    void warnIfFilesystemPartitionsExist() {
+        try {
+            for (FileStatus status : table.fileIO().listStatus(new 
Path(table.location()))) {
+                if (status.isDir() && 
!status.getPath().getName().startsWith(".")) {
+                    LOG.warn(
+                            "Format table {} has no partitions registered in 
the catalog "
+                                    + "but its location {} contains 
directories. Data written "
+                                    + "before enabling catalog-managed 
partitions (or by clients "
+                                    + "that do not register partitions) is 
invisible until the "
+                                    + "partition metadata is synced, e.g. with 
MSCK REPAIR TABLE.",
+                            table.fullName(),
+                            table.location());
+                    return;
+                }
+            }
+        } catch (IOException ignored) {
+            // Best-effort hint only; never fail or slow down the scan because 
of it.
+        }
+    }
+
+    private Map<String, String> leadingEqualityPrefix(Optional<Predicate> 
predicate) {
+        if (!predicate.isPresent()) {
+            return Collections.emptyMap();
+        }
+        return FormatTableScan.extractLeadingEqualityPartitionSpecWhenOnlyAnd(
+                table.partitionKeys(), predicate.get(), table.partitionType());
+    }
+
+    private IllegalStateException corruptPartitionSpec(@Nullable Map<String, 
String> spec) {
+        return new IllegalStateException(
+                String.format(
+                        "Catalog returned corrupt partition metadata %s for 
format table %s; "
+                                + "expected exactly the partition keys %s with 
values usable as "
+                                + "path components.",
+                        spec, table.fullName(), table.partitionKeys()));
+    }
+
+    private void warnMissingPartition(LinkedHashMap<String, String> spec, Path 
path) {
+        LOG.warn(
+                "Partition '{}' of format table {} is registered in the 
catalog but its directory "
+                        + "'{}' does not exist; treating the partition as 
empty. If the directory "
+                        + "was removed on purpose, drop the partition or 
repair the metadata, e.g. "
+                        + "with MSCK REPAIR TABLE.",
+                PartitionPathUtils.generatePartitionName(spec, false),
+                table.fullName(),
+                path);
+    }
+}
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/format/FileSystemSplitEnumerator.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/format/FileSystemSplitEnumerator.java
new file mode 100644
index 0000000000..7373e35d4d
--- /dev/null
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/format/FileSystemSplitEnumerator.java
@@ -0,0 +1,132 @@
+/*
+ * 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.paimon.table.format;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.fs.FileIO;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.manifest.PartitionEntry;
+import org.apache.paimon.partition.PartitionPredicate;
+import 
org.apache.paimon.partition.PartitionPredicate.MultiplePartitionPredicate;
+import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.table.FormatTable;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.utils.Pair;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import javax.annotation.Nullable;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Optional;
+import java.util.Set;
+
+import static 
org.apache.paimon.utils.PartitionPathUtils.searchPartSpecAndPaths;
+
+/** A {@link SplitEnumerator} whose partitions are discovered from the 
filesystem. */
+final class FileSystemSplitEnumerator extends SplitEnumerator {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(FileSystemSplitEnumerator.class);
+
+    FileSystemSplitEnumerator(FormatTable table, CoreOptions coreOptions) {
+        super(table, coreOptions);
+    }
+
+    @Override
+    List<Split> enumeratePartitions(@Nullable PartitionPredicate 
partitionFilter)
+            throws IOException {
+        List<Split> splits = new ArrayList<>();
+        FileIO fileIO = table.fileIO();
+        for (Pair<LinkedHashMap<String, String>, Path> pair : 
findPartitions(partitionFilter)) {
+            BinaryRow partitionRow = toPartitionRow(pair.getKey());
+            if (partitionFilter == null || partitionFilter.test(partitionRow)) 
{
+                splits.addAll(createSplits(fileIO, pair.getValue(), 
partitionRow));
+            }
+        }
+        return splits;
+    }
+
+    @Override
+    List<Pair<LinkedHashMap<String, String>, Path>> findPartitions(
+            @Nullable PartitionPredicate partitionFilter) {
+        LOG.debug(
+                "Find partitions for format table {}, partition filter: {}",
+                table.name(),
+                partitionFilter);
+        boolean onlyValueInPath = 
coreOptions.formatTablePartitionOnlyValueInPath();
+        if (partitionFilter instanceof MultiplePartitionPredicate) {
+            Set<BinaryRow> partitions = ((MultiplePartitionPredicate) 
partitionFilter).partitions();
+            return FormatTableScan.generatePartitions(
+                    table.partitionKeys(),
+                    table.partitionType(),
+                    table.defaultPartName(),
+                    new Path(table.location()),
+                    partitions,
+                    onlyValueInPath);
+        }
+
+        Optional<Predicate> predicate = 
FormatTableScan.extractPartitionPredicate(partitionFilter);
+        LOG.debug(
+                "Extracted predicate for format table {} partition pruning: 
{}",
+                table.name(),
+                predicate.orElse(null));
+
+        Pair<Path, Integer> scanPathAndLevel =
+                FormatTableScan.computeScanPathAndLevel(
+                        new Path(table.location()),
+                        table.partitionKeys(),
+                        predicate,
+                        table.partitionType(),
+                        onlyValueInPath);
+        return searchPartSpecAndPaths(
+                table.fileIO(),
+                scanPathAndLevel.getLeft(),
+                scanPathAndLevel.getRight(),
+                table.partitionKeys(),
+                onlyValueInPath,
+                predicate.orElse(null),
+                table.partitionType(),
+                table.defaultPartName());
+    }
+
+    @Override
+    List<PartitionEntry> listPartitionEntries() {
+        List<Pair<LinkedHashMap<String, String>, Path>> partition2Paths =
+                searchPartSpecAndPaths(
+                        table.fileIO(),
+                        new Path(table.location()),
+                        table.partitionKeys().size(),
+                        table.partitionKeys(),
+                        coreOptions.formatTablePartitionOnlyValueInPath(),
+                        null,
+                        table.partitionType(),
+                        table.defaultPartName());
+        List<PartitionEntry> partitionEntries = new ArrayList<>();
+        for (Pair<LinkedHashMap<String, String>, Path> partition2Path : 
partition2Paths) {
+            BinaryRow row = toPartitionRow(partition2Path.getKey());
+            partitionEntries.add(new PartitionEntry(row, -1L, -1L, -1L, -1L, 
-1));
+        }
+        return partitionEntries;
+    }
+}
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableScan.java 
b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableScan.java
index 2dee6fe123..72fdff6182 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableScan.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableScan.java
@@ -23,16 +23,8 @@ import org.apache.paimon.casting.CastExecutor;
 import org.apache.paimon.casting.CastExecutors;
 import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.data.BinaryString;
-import org.apache.paimon.data.GenericRow;
-import org.apache.paimon.data.serializer.InternalRowSerializer;
-import org.apache.paimon.format.csv.CsvOptions;
-import org.apache.paimon.format.json.JsonOptions;
-import org.apache.paimon.fs.FileIO;
-import org.apache.paimon.fs.FileStatus;
 import org.apache.paimon.fs.Path;
 import org.apache.paimon.manifest.PartitionEntry;
-import org.apache.paimon.options.Options;
-import org.apache.paimon.partition.Partition;
 import org.apache.paimon.partition.PartitionPredicate;
 import org.apache.paimon.partition.PartitionPredicate.AndPartitionPredicate;
 import 
org.apache.paimon.partition.PartitionPredicate.DefaultPartitionPredicate;
@@ -50,48 +42,29 @@ import org.apache.paimon.table.source.TableScan;
 import org.apache.paimon.types.DataType;
 import org.apache.paimon.types.RowType;
 import org.apache.paimon.types.VarCharType;
-import org.apache.paimon.utils.BinPacking;
 import org.apache.paimon.utils.InternalRowPartitionComputer;
 import org.apache.paimon.utils.Pair;
 import org.apache.paimon.utils.PartitionPathUtils;
 
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
 import javax.annotation.Nullable;
 
-import java.io.FileNotFoundException;
 import java.io.IOException;
 import java.util.ArrayList;
-import java.util.Arrays;
-import java.util.Collections;
-import java.util.Comparator;
 import java.util.HashMap;
-import java.util.HashSet;
 import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
 
-import static 
org.apache.paimon.format.text.HadoopCompressionUtils.isCompressed;
-import static org.apache.paimon.format.text.TextLineReader.isDefaultDelimiter;
-import static 
org.apache.paimon.utils.InternalRowPartitionComputer.convertSpecToInternalRow;
-import static 
org.apache.paimon.utils.PartitionPathUtils.searchPartSpecAndPaths;
-
 /** {@link TableScan} for {@link FormatTable}. */
 public class FormatTableScan implements InnerTableScan {
 
-    private static final Logger LOG = 
LoggerFactory.getLogger(FormatTableScan.class);
-
     final FormatTable table;
     final CoreOptions coreOptions;
     @Nullable private PartitionPredicate partitionFilter;
-    @Nullable private final FormatTablePartitionManager partitionManager;
+    private final SplitEnumerator splitEnumerator;
     @Nullable private final Integer limit;
-    private final long targetSplitSize;
-    private final long openFileCost;
-    private final FormatTable.Format format;
 
     public FormatTableScan(
             FormatTable table,
@@ -101,10 +74,7 @@ public class FormatTableScan implements InnerTableScan {
         this.coreOptions = new CoreOptions(table.options());
         this.partitionFilter = partitionFilter;
         this.limit = limit;
-        this.partitionManager = table.partitionManager();
-        this.targetSplitSize = coreOptions.splitTargetSize();
-        this.openFileCost = coreOptions.splitOpenFileCost();
-        this.format = table.format();
+        this.splitEnumerator = SplitEnumerator.create(table, coreOptions, 
table.partitionManager());
     }
 
     @Override
@@ -120,46 +90,7 @@ public class FormatTableScan implements InnerTableScan {
 
     @Override
     public List<PartitionEntry> listPartitionEntries() {
-        if (partitionManager != null) {
-            List<Partition> partitions =
-                    partitionManager.listPartitions(Collections.emptyMap(), 
null);
-            if (partitions.isEmpty()) {
-                warnIfFilesystemPartitionsExist();
-            }
-            boolean onlyValueInPath = 
coreOptions.formatTablePartitionOnlyValueInPath();
-            List<PartitionEntry> entries = new ArrayList<>(partitions.size());
-            Set<Map<String, String>> seen = new HashSet<>(partitions.size());
-            for (Partition partition : partitions) {
-                if (!seen.add(partition.spec())) {
-                    continue;
-                }
-                entries.add(
-                        new PartitionEntry(
-                                toPartitionRow(normalizeSpec(partition.spec(), 
onlyValueInPath)),
-                                partition.recordCount(),
-                                partition.fileSizeInBytes(),
-                                partition.fileCount(),
-                                partition.lastFileCreationTime(),
-                                partition.totalBuckets()));
-            }
-            return entries;
-        }
-        List<Pair<LinkedHashMap<String, String>, Path>> partition2Paths =
-                searchPartSpecAndPaths(
-                        table.fileIO(),
-                        new Path(table.location()),
-                        table.partitionKeys().size(),
-                        table.partitionKeys(),
-                        coreOptions.formatTablePartitionOnlyValueInPath(),
-                        null,
-                        table.partitionType(),
-                        table.defaultPartName());
-        List<PartitionEntry> partitionEntries = new ArrayList<>();
-        for (Pair<LinkedHashMap<String, String>, Path> partition2Path : 
partition2Paths) {
-            BinaryRow row = toPartitionRow(partition2Path.getKey());
-            partitionEntries.add(new PartitionEntry(row, -1L, -1L, -1L, -1L, 
-1));
-        }
-        return partitionEntries;
+        return splitEnumerator.listPartitionEntries();
     }
 
     @Override
@@ -172,10 +103,7 @@ public class FormatTableScan implements InnerTableScan {
     }
 
     BinaryRow toPartitionRow(LinkedHashMap<String, String> partitionSpec) {
-        RowType partitionType = table.partitionType();
-        GenericRow row =
-                convertSpecToInternalRow(partitionSpec, partitionType, 
table.defaultPartName());
-        return new InternalRowSerializer(partitionType).toBinaryRow(row);
+        return splitEnumerator.toPartitionRow(partitionSpec);
     }
 
     private class FormatTableScanPlan implements Plan {
@@ -183,39 +111,7 @@ public class FormatTableScan implements InnerTableScan {
         public List<Split> splits() {
             List<Split> splits = new ArrayList<>();
             try {
-                FileIO fileIO = table.fileIO();
-                if (!table.partitionKeys().isEmpty()) {
-                    for (Pair<LinkedHashMap<String, String>, Path> pair : 
findPartitions()) {
-                        LinkedHashMap<String, String> partitionSpec = 
pair.getKey();
-                        BinaryRow partitionRow = toPartitionRow(partitionSpec);
-                        if (partitionFilter == null || 
partitionFilter.test(partitionRow)) {
-                            try {
-                                splits.addAll(createSplits(fileIO, 
pair.getValue(), partitionRow));
-                            } catch (FileNotFoundException e) {
-                                if (partitionManager == null) {
-                                    throw e;
-                                }
-                                // A registered partition without a directory 
reads as empty,
-                                // matching Hive (e.g. ADD PARTITION before 
the first INSERT).
-                                // Warn so a directory removed behind the 
catalog's back stays
-                                // discoverable.
-                                LOG.warn(
-                                        "Partition '{}' of format table {} is 
registered in "
-                                                + "the catalog but its 
directory '{}' does not "
-                                                + "exist; treating the 
partition as empty. If the "
-                                                + "directory was removed on 
purpose, drop the "
-                                                + "partition or repair the 
metadata, e.g. with "
-                                                + "MSCK REPAIR TABLE.",
-                                        
PartitionPathUtils.generatePartitionName(
-                                                partitionSpec, false),
-                                        table.fullName(),
-                                        pair.getValue());
-                            }
-                        }
-                    }
-                } else {
-                    splits.addAll(createSplits(fileIO, new 
Path(table.location()), null));
-                }
+                splits.addAll(splitEnumerator.enumerate(partitionFilter));
                 // Keep all splits for a positive limit because 
FormatDataSplit has no row count.
                 if (limit != null && limit <= 0) {
                     return new ArrayList<>();
@@ -228,156 +124,7 @@ public class FormatTableScan implements InnerTableScan {
     }
 
     List<Pair<LinkedHashMap<String, String>, Path>> findPartitions() {
-        if (partitionManager != null) {
-            Optional<Predicate> extracted = 
extractPartitionPredicate(partitionFilter);
-            Map<String, String> prefix = leadingEqualityPrefix(extracted);
-            // The whole predicate goes to the manager as a pushdown hint, 
together with the
-            // prefix; the catalog may return a superset.
-            Predicate catalogFilter = extracted.orElse(null);
-            List<Partition> partitions = 
partitionManager.listPartitions(prefix, catalogFilter);
-            if (partitions.isEmpty() && prefix.isEmpty() && catalogFilter == 
null) {
-                warnIfFilesystemPartitionsExist();
-            }
-            // The prefix and filter are coarse pre-filters; the full 
predicate is applied per
-            // partition in the plan, as it is for filesystem discovery.
-            return toSpecsAndPaths(partitions, 
coreOptions.formatTablePartitionOnlyValueInPath());
-        }
-        LOG.debug(
-                "Find partitions for format table {}, partition filter: {}",
-                table.name(),
-                partitionFilter);
-        boolean onlyValueInPath = 
coreOptions.formatTablePartitionOnlyValueInPath();
-        if (partitionFilter instanceof MultiplePartitionPredicate) {
-            // generate partitions directly
-            Set<BinaryRow> partitions = ((MultiplePartitionPredicate) 
partitionFilter).partitions();
-            return generatePartitions(
-                    table.partitionKeys(),
-                    table.partitionType(),
-                    table.defaultPartName(),
-                    new Path(table.location()),
-                    partitions,
-                    onlyValueInPath);
-        } else {
-            // search paths with partition filter optimization
-            // This will prune partition directories early during traversal,
-            // which is especially important for cloud storage like OSS/S3
-            Optional<Predicate> predicate = 
extractPartitionPredicate(partitionFilter);
-            LOG.debug(
-                    "Extracted predicate for format table {} partition 
pruning: {}",
-                    table.name(),
-                    predicate.orElse(null));
-
-            Pair<Path, Integer> scanPathAndLevel =
-                    computeScanPathAndLevel(
-                            new Path(table.location()),
-                            table.partitionKeys(),
-                            predicate,
-                            table.partitionType(),
-                            onlyValueInPath);
-            return searchPartSpecAndPaths(
-                    table.fileIO(),
-                    scanPathAndLevel.getLeft(),
-                    scanPathAndLevel.getRight(),
-                    table.partitionKeys(),
-                    onlyValueInPath,
-                    predicate.orElse(null),
-                    table.partitionType(),
-                    table.defaultPartName());
-        }
-    }
-
-    /**
-     * Turn the partitions a catalog reports into the specs and directories to 
read. Catalog
-     * metadata is not trusted for path construction: a spec that cannot form 
a partition directory
-     * of this table is rejected rather than resolved to some other directory.
-     */
-    private List<Pair<LinkedHashMap<String, String>, Path>> toSpecsAndPaths(
-            List<Partition> partitions, boolean onlyValueInPath) {
-        List<Pair<LinkedHashMap<String, String>, Path>> result = new 
ArrayList<>(partitions.size());
-        Path tablePath = new Path(table.location());
-        // Do not trust the catalog to be duplicate-free: a repeated spec 
would double every split
-        // of that partition and silently duplicate query results.
-        Set<String> seenPartitionPaths = new HashSet<>(partitions.size());
-        for (Partition partition : partitions) {
-            LinkedHashMap<String, String> spec = 
normalizeSpec(partition.spec(), onlyValueInPath);
-            String partitionPath =
-                    PartitionPathUtils.generatePartitionPathUtil(spec, 
onlyValueInPath);
-            if (seenPartitionPaths.add(partitionPath)) {
-                result.add(Pair.of(spec, new Path(tablePath, partitionPath)));
-            }
-        }
-        return result;
-    }
-
-    /** Order a catalog spec by the table partition keys and check it can form 
a directory. */
-    private LinkedHashMap<String, String> normalizeSpec(
-            @Nullable Map<String, String> spec, boolean onlyValueInPath) {
-        List<String> partitionKeys = table.partitionKeys();
-        if (spec == null
-                || spec.size() != partitionKeys.size()
-                || !spec.keySet().containsAll(partitionKeys)) {
-            throw corruptPartitionSpec(spec);
-        }
-        LinkedHashMap<String, String> normalized = new LinkedHashMap<>();
-        for (String partitionKey : partitionKeys) {
-            String value = spec.get(partitionKey);
-            // In a value-only layout, '.' and '..' are complete path 
components and would resolve
-            // to a directory outside the table.
-            try {
-                PartitionPathUtils.validatePartitionValueForPath(value, 
onlyValueInPath);
-            } catch (IllegalArgumentException e) {
-                throw corruptPartitionSpec(spec);
-            }
-            normalized.put(partitionKey, value);
-        }
-        return normalized;
-    }
-
-    private IllegalStateException corruptPartitionSpec(@Nullable Map<String, 
String> spec) {
-        return new IllegalStateException(
-                String.format(
-                        "Catalog returned corrupt partition metadata %s for 
format table %s; "
-                                + "expected exactly the partition keys %s with 
values usable as "
-                                + "path components.",
-                        spec, table.fullName(), table.partitionKeys()));
-    }
-
-    /**
-     * The leading equality prefix of the partition predicate, in 
partition-key order, pushed down
-     * to the partition catalog. Empty when there is no predicate or it does 
not start with
-     * equalities on the leading partition keys.
-     */
-    private Map<String, String> leadingEqualityPrefix(Optional<Predicate> 
predicate) {
-        if (!predicate.isPresent()) {
-            return Collections.emptyMap();
-        }
-        return extractLeadingEqualityPartitionSpecWhenOnlyAnd(
-                table.partitionKeys(), predicate.get(), table.partitionType());
-    }
-
-    /**
-     * Warn when the catalog knows no partitions but the table directory 
contains subdirectories:
-     * typically a table that predates catalog-managed partitions (or was 
written by a client that
-     * does not register partitions) and needs a metadata sync before its data 
becomes visible.
-     */
-    private void warnIfFilesystemPartitionsExist() {
-        try {
-            for (FileStatus status : table.fileIO().listStatus(new 
Path(table.location()))) {
-                if (status.isDir() && 
!status.getPath().getName().startsWith(".")) {
-                    LOG.warn(
-                            "Format table {} has no partitions registered in 
the catalog "
-                                    + "but its location {} contains 
directories. Data written "
-                                    + "before enabling catalog-managed 
partitions (or by clients "
-                                    + "that do not register partitions) is 
invisible until the "
-                                    + "partition metadata is synced, e.g. with 
MSCK REPAIR TABLE.",
-                            table.fullName(),
-                            table.location());
-                    return;
-                }
-            }
-        } catch (IOException ignored) {
-            // Best-effort hint only; never fail or slow down the scan because 
of it.
-        }
+        return splitEnumerator.findPartitions(partitionFilter);
     }
 
     protected static List<Pair<LinkedHashMap<String, String>, Path>> 
generatePartitions(
@@ -426,7 +173,7 @@ public class FormatTableScan implements InnerTableScan {
                 Optional<Predicate> childPredicate = 
extractPartitionPredicate(child);
                 childPredicate.ifPresent(predicates::add);
                 // Skip children that can't be expressed as Predicate (e.g. 
Multiple);
-                // they are still applied by partitionFilter.test() in plan().
+                // they are still applied before listing the partition files.
             }
             return predicates.isEmpty()
                     ? Optional.empty()
@@ -461,66 +208,6 @@ public class FormatTableScan implements InnerTableScan {
         return Pair.of(scanPath, level);
     }
 
-    private List<Split> createSplits(FileIO fileIO, Path path, BinaryRow 
partition)
-            throws IOException {
-        List<FormatDataSplit.FileMeta> segments = new ArrayList<>();
-        FileStatus[] files = fileIO.listFiles(path, true);
-        Arrays.sort(files, Comparator.comparing(file -> 
file.getPath().toString()));
-        for (FileStatus file : files) {
-            if (isDataFileName(file.getPath().getName())) {
-                segments.addAll(toSegments(file));
-            }
-        }
-
-        List<Split> splits = new ArrayList<>();
-        for (List<FormatDataSplit.FileMeta> bin :
-                BinPacking.packForOrdered(
-                        segments,
-                        file -> Math.max(file.readSize(), openFileCost),
-                        targetSplitSize)) {
-            splits.add(new FormatDataSplit(bin, partition));
-        }
-        return splits;
-    }
-
-    private List<FormatDataSplit.FileMeta> toSegments(FileStatus file) {
-        if (!preferToSplitFile(file)) {
-            return Collections.singletonList(
-                    new FormatDataSplit.FileMeta(file.getPath(), 
file.getLen()));
-        }
-        List<FormatDataSplit.FileMeta> segments = new ArrayList<>();
-        long remainingBytes = file.getLen();
-        long currentStart = 0;
-
-        while (remainingBytes > 0) {
-            long splitSize = Math.min(targetSplitSize, remainingBytes);
-            segments.add(
-                    new FormatDataSplit.FileMeta(
-                            file.getPath(), file.getLen(), currentStart, 
splitSize));
-            currentStart += splitSize;
-            remainingBytes -= splitSize;
-        }
-        return segments;
-    }
-
-    private boolean preferToSplitFile(FileStatus file) {
-        if (file.getLen() <= targetSplitSize) {
-            return false;
-        }
-
-        Options options = coreOptions.toConfiguration();
-        switch (format) {
-            case CSV:
-                return !isCompressed(file.getPath())
-                        && 
isDefaultDelimiter(options.get(CsvOptions.LINE_DELIMITER));
-            case JSON:
-                return !isCompressed(file.getPath())
-                        && 
isDefaultDelimiter(options.get(JsonOptions.LINE_DELIMITER));
-            default:
-                return false;
-        }
-    }
-
     public static Map<String, String> 
extractLeadingEqualityPartitionSpecWhenOnlyAnd(
             List<String> partitionKeys, Predicate predicate, RowType 
partitionType) {
         List<Predicate> predicates = PredicateBuilder.splitAnd(predicate);
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/format/SplitEnumerator.java 
b/paimon-core/src/main/java/org/apache/paimon/table/format/SplitEnumerator.java
new file mode 100644
index 0000000000..4bfea41a44
--- /dev/null
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/format/SplitEnumerator.java
@@ -0,0 +1,168 @@
+/*
+ * 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.paimon.table.format;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.serializer.InternalRowSerializer;
+import org.apache.paimon.format.csv.CsvOptions;
+import org.apache.paimon.format.json.JsonOptions;
+import org.apache.paimon.fs.FileIO;
+import org.apache.paimon.fs.FileStatus;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.manifest.PartitionEntry;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.partition.PartitionPredicate;
+import org.apache.paimon.table.FormatTable;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.BinPacking;
+import org.apache.paimon.utils.Pair;
+
+import javax.annotation.Nullable;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.LinkedHashMap;
+import java.util.List;
+
+import static 
org.apache.paimon.format.text.HadoopCompressionUtils.isCompressed;
+import static org.apache.paimon.format.text.TextLineReader.isDefaultDelimiter;
+import static 
org.apache.paimon.utils.InternalRowPartitionComputer.convertSpecToInternalRow;
+
+/**
+ * Enumerates {@link FormatDataSplit}s for a {@link FormatTable}.
+ *
+ * <p>Partitions come either from the catalog ({@link CatalogSplitEnumerator}) 
or from the
+ * filesystem ({@link FileSystemSplitEnumerator}); the source is chosen once 
by {@link #create} so
+ * the enumeration paths stay free of {@code partitionManager != null} 
branching.
+ */
+abstract class SplitEnumerator {
+
+    protected final FormatTable table;
+    protected final CoreOptions coreOptions;
+    protected final long targetSplitSize;
+    protected final long openFileCost;
+    protected final FormatTable.Format format;
+
+    SplitEnumerator(FormatTable table, CoreOptions coreOptions) {
+        this.table = table;
+        this.coreOptions = coreOptions;
+        this.targetSplitSize = coreOptions.splitTargetSize();
+        this.openFileCost = coreOptions.splitOpenFileCost();
+        this.format = table.format();
+    }
+
+    static SplitEnumerator create(
+            FormatTable table,
+            CoreOptions coreOptions,
+            @Nullable FormatTablePartitionManager partitionManager) {
+        if (partitionManager != null) {
+            return new CatalogSplitEnumerator(table, coreOptions, 
partitionManager);
+        }
+        return new FileSystemSplitEnumerator(table, coreOptions);
+    }
+
+    final List<Split> enumerate(@Nullable PartitionPredicate partitionFilter) 
throws IOException {
+        if (table.partitionKeys().isEmpty()) {
+            return createSplits(table.fileIO(), new Path(table.location()), 
null);
+        }
+        return enumeratePartitions(partitionFilter);
+    }
+
+    /** Enumerate splits for a partitioned table; partitions come from the 
concrete source. */
+    abstract List<Split> enumeratePartitions(@Nullable PartitionPredicate 
partitionFilter)
+            throws IOException;
+
+    abstract List<Pair<LinkedHashMap<String, String>, Path>> findPartitions(
+            @Nullable PartitionPredicate partitionFilter);
+
+    abstract List<PartitionEntry> listPartitionEntries();
+
+    BinaryRow toPartitionRow(LinkedHashMap<String, String> partitionSpec) {
+        RowType partitionType = table.partitionType();
+        GenericRow row =
+                convertSpecToInternalRow(partitionSpec, partitionType, 
table.defaultPartName());
+        return new InternalRowSerializer(partitionType).toBinaryRow(row);
+    }
+
+    List<Split> createSplits(FileIO fileIO, Path path, @Nullable BinaryRow 
partition)
+            throws IOException {
+        List<FormatDataSplit.FileMeta> segments = new ArrayList<>();
+        FileStatus[] files = fileIO.listFiles(path, true);
+        Arrays.sort(files, Comparator.comparing(file -> 
file.getPath().toString()));
+        for (FileStatus file : files) {
+            if (FormatTableScan.isDataFileName(file.getPath().getName())) {
+                segments.addAll(toSegments(file));
+            }
+        }
+
+        List<Split> splits = new ArrayList<>();
+        for (List<FormatDataSplit.FileMeta> bin :
+                BinPacking.packForOrdered(
+                        segments,
+                        file -> Math.max(file.readSize(), openFileCost),
+                        targetSplitSize)) {
+            splits.add(new FormatDataSplit(bin, partition));
+        }
+        return splits;
+    }
+
+    private List<FormatDataSplit.FileMeta> toSegments(FileStatus file) {
+        if (!preferToSplitFile(file)) {
+            return Collections.singletonList(
+                    new FormatDataSplit.FileMeta(file.getPath(), 
file.getLen()));
+        }
+        List<FormatDataSplit.FileMeta> segments = new ArrayList<>();
+        long remainingBytes = file.getLen();
+        long currentStart = 0;
+
+        while (remainingBytes > 0) {
+            long splitSize = Math.min(targetSplitSize, remainingBytes);
+            segments.add(
+                    new FormatDataSplit.FileMeta(
+                            file.getPath(), file.getLen(), currentStart, 
splitSize));
+            currentStart += splitSize;
+            remainingBytes -= splitSize;
+        }
+        return segments;
+    }
+
+    private boolean preferToSplitFile(FileStatus file) {
+        if (file.getLen() <= targetSplitSize) {
+            return false;
+        }
+
+        Options options = coreOptions.toConfiguration();
+        switch (format) {
+            case CSV:
+                return !isCompressed(file.getPath())
+                        && 
isDefaultDelimiter(options.get(CsvOptions.LINE_DELIMITER));
+            case JSON:
+                return !isCompressed(file.getPath())
+                        && 
isDefaultDelimiter(options.get(JsonOptions.LINE_DELIMITER));
+            default:
+                return false;
+        }
+    }
+}
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/format/CatalogManagedPartitionScanTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/format/CatalogManagedPartitionScanTest.java
index de2127f493..be7d5e4584 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/format/CatalogManagedPartitionScanTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/format/CatalogManagedPartitionScanTest.java
@@ -53,9 +53,14 @@ import java.util.Collections;
 import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
 import java.util.stream.Collectors;
 
 import static 
org.apache.paimon.CoreOptions.FORMAT_TABLE_PARTITION_ONLY_VALUE_IN_PATH;
+import static org.apache.paimon.CoreOptions.FORMAT_TABLE_SCAN_LIST_PARALLELISM;
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 import static org.mockito.ArgumentMatchers.any;
@@ -506,4 +511,281 @@ class CatalogManagedPartitionScanTest {
             return super.listFiles(path, recursive);
         }
     }
+
+    @Test
+    void testCatalogPartitionListingRunsInParallelAndPreservesOrder() throws 
Exception {
+        ParallelTrackingLocalFileIO fileIO = new ParallelTrackingLocalFileIO();
+        Path tablePath = new Path(tempDir.toUri());
+        List<Partition> partitions = new ArrayList<>();
+        List<Path> expected = new ArrayList<>();
+        for (int m = 10; m <= 15; m++) {
+            partitions.add(partition("2025", Integer.toString(m)));
+            expected.add(writeDataFile(fileIO, tablePath, "year=2025/month=" + 
m));
+        }
+
+        List<Path> parallel =
+                plannedFiles(
+                        new FormatTableScan(
+                                        stringPartitionTable(
+                                                fileIO, tablePath, 
recordingCatalog(partitions), 8),
+                                        null,
+                                        null)
+                                .plan()
+                                .splits());
+
+        assertThat(parallel).containsExactlyElementsOf(expected);
+        assertThat(fileIO.maxConcurrentListings()).isGreaterThan(1);
+    }
+
+    @Test
+    void testParallelListingSkipsMissingCatalogRegisteredPartition() throws 
Exception {
+        InjectingLocalFileIO fileIO = new InjectingLocalFileIO();
+        Path tablePath = new Path(tempDir.toUri());
+        List<Partition> partitions =
+                Arrays.asList(
+                        partition("2025", "10"), partition("2025", "11"), 
partition("2025", "12"));
+        Path oct = writeDataFile(fileIO, tablePath, "year=2025/month=10");
+        Path dec = writeDataFile(fileIO, tablePath, "year=2025/month=12");
+        fileIO.failListFilesContaining("month=11", new 
FileNotFoundException("missing"));
+
+        List<Path> files =
+                plannedFiles(
+                        new FormatTableScan(
+                                        stringPartitionTable(
+                                                fileIO, tablePath, 
recordingCatalog(partitions), 4),
+                                        null,
+                                        null)
+                                .plan()
+                                .splits());
+
+        assertThat(files).containsExactlyInAnyOrder(oct, dec);
+    }
+
+    @Test
+    void testCatalogPartitionListingIoErrorFailsWholeScan() throws Exception {
+        InjectingLocalFileIO fileIO = new InjectingLocalFileIO();
+        Path tablePath = new Path(tempDir.toUri());
+        List<Partition> partitions =
+                Arrays.asList(partition("2025", "10"), partition("2025", 
"11"));
+        writeDataFile(fileIO, tablePath, "year=2025/month=10");
+        writeDataFile(fileIO, tablePath, "year=2025/month=11");
+        fileIO.failListFilesContaining("month=11", new IOException("boom"));
+
+        FormatTable table =
+                stringPartitionTable(fileIO, tablePath, 
recordingCatalog(partitions), 4);
+        assertThatThrownBy(() -> new FormatTableScan(table, null, 
null).plan().splits())
+                .isInstanceOf(RuntimeException.class)
+                .hasRootCauseInstanceOf(IOException.class)
+                .hasRootCauseMessage("boom");
+    }
+
+    @Test
+    void testFilesystemDiscoveredMissingDirectoryStillFailsScan() throws 
Exception {
+        InjectingLocalFileIO fileIO = new InjectingLocalFileIO();
+        Path tablePath = new Path(tempDir.toUri());
+        writeDataFile(fileIO, tablePath, "year=2025/month=10");
+        writeDataFile(fileIO, tablePath, "year=2025/month=11");
+        fileIO.failListFilesContaining("month=11", new 
FileNotFoundException("missing"));
+
+        FormatTable table = createStringPartitionTable(fileIO, tablePath, 
null);
+        assertThatThrownBy(() -> new FormatTableScan(table, null, 
null).plan().splits())
+                .isInstanceOf(RuntimeException.class)
+                .hasRootCauseInstanceOf(FileNotFoundException.class);
+    }
+
+    @Test
+    void testListParallelismDoesNotChangeFilesystemDiscoveredScanning() throws 
Exception {
+        SerialTrackingLocalFileIO fileIO = new SerialTrackingLocalFileIO();
+        Path tablePath = new Path(tempDir.toUri());
+        List<Path> expected = new ArrayList<>();
+        for (int month = 10; month <= 12; month++) {
+            expected.add(writeDataFile(fileIO, tablePath, "year=2025/month=" + 
month));
+        }
+
+        FormatTable table = stringPartitionTable(fileIO, tablePath, null, 8);
+        List<Path> files = plannedFiles(new FormatTableScan(table, null, 
null).plan().splits());
+
+        assertThat(files).containsExactlyInAnyOrderElementsOf(expected);
+        assertThat(fileIO.maxConcurrentListings()).isOne();
+    }
+
+    @Test
+    void testCatalogListingEstablishesFilesystemOnCallerThread() throws 
Exception {
+        CallerBoundLocalFileIO fileIO = new CallerBoundLocalFileIO();
+        Path tablePath = new Path(tempDir.toUri());
+        List<Partition> partitions = new ArrayList<>();
+        for (int month = 10; month <= 13; month++) {
+            partitions.add(partition("2025", Integer.toString(month)));
+            writeDataFile(fileIO, tablePath, "year=2025/month=" + month);
+        }
+
+        fileIO.resetBinding();
+        CallerBoundLocalFileIO.CALLER_USER.set("caller-user");
+        try {
+            new FormatTableScan(
+                            stringPartitionTable(
+                                    fileIO, tablePath, 
recordingCatalog(partitions), 4),
+                            null,
+                            null)
+                    .plan()
+                    .splits();
+        } finally {
+            CallerBoundLocalFileIO.CALLER_USER.remove();
+        }
+
+        assertThat(fileIO.boundUser()).isEqualTo("caller-user");
+        assertThat(fileIO.boundOnListingWorker()).isFalse();
+    }
+
+    private FormatTable stringPartitionTable(
+            LocalFileIO fileIO,
+            Path tablePath,
+            @Nullable FormatTablePartitionManager partitionManager,
+            int parallelism) {
+        RowType rowType =
+                RowType.builder()
+                        .field("year", DataTypes.STRING())
+                        .field("month", DataTypes.STRING())
+                        .field("id", DataTypes.INT())
+                        .build();
+        Map<String, String> options = new LinkedHashMap<>();
+        options.put(FORMAT_TABLE_PARTITION_ONLY_VALUE_IN_PATH.key(), "false");
+        options.put(FORMAT_TABLE_SCAN_LIST_PARALLELISM.key(), 
Integer.toString(parallelism));
+        return FormatTable.builder()
+                .fileIO(fileIO)
+                .identifier(IDENTIFIER)
+                .rowType(rowType)
+                .partitionKeys(Arrays.asList("year", "month"))
+                .location(tablePath.toString())
+                .format(FormatTable.Format.CSV)
+                .options(options)
+                .partitionManager(partitionManager)
+                .build();
+    }
+
+    private static class ParallelTrackingLocalFileIO extends LocalFileIO {
+
+        private final CountDownLatch concurrentListings = new 
CountDownLatch(2);
+        private final AtomicInteger activeListings = new AtomicInteger();
+        private final AtomicInteger maxConcurrentListings = new 
AtomicInteger();
+
+        @Override
+        public FileStatus[] listFiles(Path path, boolean recursive) throws 
IOException {
+            int active = activeListings.incrementAndGet();
+            maxConcurrentListings.updateAndGet(current -> Math.max(current, 
active));
+            concurrentListings.countDown();
+            try {
+                if (!concurrentListings.await(10, TimeUnit.SECONDS)) {
+                    throw new IOException("Partition file listings did not run 
concurrently");
+                }
+                return super.listFiles(path, recursive);
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+                throw new IOException("Interrupted while waiting for 
concurrent listings", e);
+            } finally {
+                activeListings.decrementAndGet();
+            }
+        }
+
+        private int maxConcurrentListings() {
+            return maxConcurrentListings.get();
+        }
+    }
+
+    private static class SerialTrackingLocalFileIO extends LocalFileIO {
+
+        private final AtomicInteger activeListings = new AtomicInteger();
+        private final AtomicInteger maxConcurrentListings = new 
AtomicInteger();
+
+        @Override
+        public FileStatus[] listFiles(Path path, boolean recursive) throws 
IOException {
+            int active = activeListings.incrementAndGet();
+            maxConcurrentListings.updateAndGet(current -> Math.max(current, 
active));
+            try {
+                Thread.sleep(50);
+                return super.listFiles(path, recursive);
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+                throw new IOException("Interrupted while tracking file 
listings", e);
+            } finally {
+                activeListings.decrementAndGet();
+            }
+        }
+
+        private int maxConcurrentListings() {
+            return maxConcurrentListings.get();
+        }
+    }
+
+    private static class InjectingLocalFileIO extends LocalFileIO {
+
+        private final Map<String, IOException> listFailures = new 
LinkedHashMap<>();
+
+        void failListFilesContaining(String marker, IOException failure) {
+            listFailures.put(marker, failure);
+        }
+
+        @Override
+        public FileStatus[] listFiles(Path path, boolean recursive) throws 
IOException {
+            for (Map.Entry<String, IOException> entry : 
listFailures.entrySet()) {
+                if (path.toString().contains(entry.getKey())) {
+                    throw entry.getValue();
+                }
+            }
+            return super.listFiles(path, recursive);
+        }
+    }
+
+    /** Binds to the thread and thread-local user that first touches the 
filesystem. */
+    private static class CallerBoundLocalFileIO extends LocalFileIO {
+
+        static final ThreadLocal<String> CALLER_USER = new ThreadLocal<>();
+
+        private final AtomicReference<Thread> boundThread = new 
AtomicReference<>();
+        private final AtomicReference<String> boundUser = new 
AtomicReference<>();
+
+        private void bindOnce() {
+            if (boundThread.compareAndSet(null, Thread.currentThread())) {
+                boundUser.set(CALLER_USER.get());
+            }
+        }
+
+        void resetBinding() {
+            boundThread.set(null);
+            boundUser.set(null);
+        }
+
+        String boundUser() {
+            return boundUser.get();
+        }
+
+        boolean boundOnListingWorker() {
+            Thread thread = boundThread.get();
+            return thread != null && 
thread.getName().contains("FORMAT-TABLE-LIST");
+        }
+
+        @Override
+        public boolean exists(Path path) throws IOException {
+            bindOnce();
+            return super.exists(path);
+        }
+
+        @Override
+        public FileStatus getFileStatus(Path path) throws IOException {
+            bindOnce();
+            return super.getFileStatus(path);
+        }
+
+        @Override
+        public FileStatus[] listStatus(Path path) throws IOException {
+            bindOnce();
+            return super.listStatus(path);
+        }
+
+        @Override
+        public FileStatus[] listFiles(Path path, boolean recursive) throws 
IOException {
+            bindOnce();
+            return super.listFiles(path, recursive);
+        }
+    }
 }

Reply via email to