Gabriel39 commented on code in PR #66247:
URL: https://github.com/apache/doris/pull/66247#discussion_r3690587646


##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java:
##########
@@ -316,20 +317,60 @@ Table resolveTable(PaimonTableHandle paimonHandle) {
     Table resolveScanTable(PaimonTableHandle paimonHandle) {
         Table table = resolveTable(paimonHandle);
         Map<String, String> scanOptions = paimonHandle.getScanOptions();
+        Table finalTable = table;
         if (scanOptions != null && !scanOptions.isEmpty()) {
             if (PaimonScanParams.isOptionsPin(scanOptions)) {
                 // An @options pin owns the whole scan-startup state: 
applyOptions strips the internal
                 // markers and nulls out the absent members of paimon's 
inherited read-state family, so a
                 // scan.mode / tag persisted on the base table cannot leak 
into this relation's read.
-                return PaimonScanParams.applyOptions(table, scanOptions);
+                finalTable = PaimonScanParams.applyOptions(table, scanOptions);
+            } else {
+                // FIX-INCR-SCAN-RESET: for an @incr read, reapply legacy's 
null reset of
+                // scan.snapshot-id/scan.mode here (the single Table.copy 
chokepoint shared by both the
+                // native/JNI scan path and the JNI serialized-table path) so 
a stale persisted pin on the
+                // base table cannot hijack incremental-between. 
Non-incremental pins pass through unchanged.
+                finalTable = 
table.copy(PaimonIncrementalScanParams.applyResetsIfIncremental(scanOptions));
+            }
+        }
+        finalTable = runtimeSafeTable(finalTable);
+        // This is the last common boundary before planning and serialization. 
Normalize and
+        // validate only after relation > catalog > physical precedence is 
established.
+        PaimonReaderOptions.validateEffectiveTable(finalTable);
+        validateHiddenSystemDataTable(paimonHandle, scanOptions);
+        return finalTable;
+    }
+
+    private Table runtimeSafeTable(Table table) {
+        Map<String, String> runtimeOptions = 
PaimonReaderOptions.runtimeSafeCopyOptions(
+                table, Collections.emptyMap());
+        // The cached catalog handle remains hardware-neutral; only the 
query-local planning copy
+        // receives a CPU-local cap before it can resize Paimon's JVM-wide 
manifest executor.
+        return runtimeOptions.isEmpty() ? table : table.copy(runtimeOptions);
+    }
+
+    private void validateHiddenSystemDataTable(PaimonTableHandle handle, 
Map<String, String> scanOptions) {
+        if (!handle.isSystemTable()) {
+            return;
+        }
+        try {
+            Table dataTable = handle.getSystemTableSource();
+            if (dataTable == null) {
+                dataTable = handle.getSysBaseTable();
+            }
+            if (dataTable == null) {
+                dataTable = catalogOps.getTable(
+                        Identifier.create(handle.getDatabaseName(), 
handle.getTableName()));
             }
-            // FIX-INCR-SCAN-RESET: for an @incr read, reapply legacy's null 
reset of
-            // scan.snapshot-id/scan.mode here (the single Table.copy 
chokepoint shared by both the
-            // native/JNI scan path and the JNI serialized-table path) so a 
stale persisted pin on the
-            // base table cannot hijack incremental-between. Non-incremental 
pins pass through unchanged.
-            return 
table.copy(PaimonIncrementalScanParams.applyResetsIfIncremental(scanOptions));
+            if (PaimonScanParams.isOptionsPin(scanOptions)) {
+                // Read-only system wrappers plan manifests through their 
hidden data table, so the
+                // same relation copy must establish precedence on both 
visible and hidden handles.
+                PaimonScanParams.applyOptions(dataTable, scanOptions);

Review Comment:
   Valid. The current head rebuilds system wrappers from the normalized hidden 
source instead of discarding the capped copy.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonConnectorMetadata.java:
##########
@@ -1444,7 +1454,10 @@ public Optional<ConnectorTableStatistics> 
getTableStatistics(
         PaimonTableHandle paimonHandle = (PaimonTableHandle) handle;
         long rowCount;
         try {
-            rowCount = catalogOps.rowCount(resolveTable(paimonHandle));
+            Table table = resolveTable(paimonHandle);
+            PaimonReaderOptions.validateEffectiveTable(table);

Review Comment:
   Valid. Statistics now use the runtime-safe execution table before validation 
and row-count planning.



##########
fe/be-java-extensions/paimon-scanner/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java:
##########
@@ -575,12 +670,31 @@ static Optional<Long> parseDataSizeBytes(String value) {
     private void initTable() {
         Preconditions.checkState(params.containsKey("serialized_table"));
         table = PaimonUtils.deserialize(params.get("serialized_table"));
+        
validateSerializedReadBatchSize(table.options().get(CoreOptions.READ_BATCH_SIZE.key()));

Review Comment:
   Valid. The rolling-upgrade scanner guard recursively validates hidden 
readable children.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java:
##########
@@ -316,20 +317,60 @@ Table resolveTable(PaimonTableHandle paimonHandle) {
     Table resolveScanTable(PaimonTableHandle paimonHandle) {
         Table table = resolveTable(paimonHandle);
         Map<String, String> scanOptions = paimonHandle.getScanOptions();
+        Table finalTable = table;
         if (scanOptions != null && !scanOptions.isEmpty()) {
             if (PaimonScanParams.isOptionsPin(scanOptions)) {
                 // An @options pin owns the whole scan-startup state: 
applyOptions strips the internal
                 // markers and nulls out the absent members of paimon's 
inherited read-state family, so a
                 // scan.mode / tag persisted on the base table cannot leak 
into this relation's read.
-                return PaimonScanParams.applyOptions(table, scanOptions);
+                finalTable = PaimonScanParams.applyOptions(table, scanOptions);
+            } else {
+                // FIX-INCR-SCAN-RESET: for an @incr read, reapply legacy's 
null reset of
+                // scan.snapshot-id/scan.mode here (the single Table.copy 
chokepoint shared by both the
+                // native/JNI scan path and the JNI serialized-table path) so 
a stale persisted pin on the
+                // base table cannot hijack incremental-between. 
Non-incremental pins pass through unchanged.
+                finalTable = 
table.copy(PaimonIncrementalScanParams.applyResetsIfIncremental(scanOptions));
+            }
+        }
+        finalTable = runtimeSafeTable(finalTable);

Review Comment:
   Valid. Hidden fallback and delegated planners are normalized, and deferred 
system planning carries a smaller-BE cap.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonCatalogOps.java:
##########
@@ -267,21 +268,24 @@ public List<String> listTables(String databaseName) 
throws Catalog.DatabaseNotEx
          * #65955: overlay the catalog-level {@code paimon.table-option.*} 
defaults onto the loaded
          * table, exactly where legacy {@code 
PaimonExternalCatalog.getPaimonTable} did — this is the
          * connector's only {@code Catalog.getTable} call, so branch, 
time-travel and system tables all
-         * inherit the defaults. Options the table sets itself win (see {@link 
PaimonTableOptions#forCopy}).
+         * inherit the configured reader policy. Catalog options override 
physical table values;
+         * relation-scoped options can still override them for one scan.
          */
         @Override
         public Table getTable(Identifier identifier) throws 
Catalog.TableNotExistException {
             Table table = catalog.getTable(identifier);
-            if (tableOptions.isEmpty()) {
-                return table;
-            }
-            Map<String, String> optionsForCopy = 
PaimonTableOptions.forCopy(tableOptions, table.options());
+            Map<String, String> optionsForCopy = 
PaimonTableOptions.forCopy(tableOptions);
+            // Relation options are applied after this cached handle is 
returned. Defer final
+            // validation so a safe relation value can override an unsafe 
physical value.
             return optionsForCopy.isEmpty() ? table : 
table.copy(optionsForCopy);
         }
 
         @Override
-        public List<Partition> listPartitions(Identifier identifier) throws 
Catalog.TableNotExistException {
-            return catalog.listPartitions(identifier);
+        public List<Partition> listPartitions(Identifier identifier, Table 
table)
+                throws Catalog.TableNotExistException {
+            // The supplied handle already contains catalog and relation 
policy. Reloading by identifier
+            // would discard those copies before manifest enumeration reaches 
the final scan guard.
+            return CatalogUtils.listPartitionsFromFileSystem(table);

Review Comment:
   Valid. REST-owned partition results remain authoritative; the supplied 
effective table is used only for filesystem fallback.



##########
fe/be-java-extensions/paimon-scanner/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java:
##########
@@ -575,12 +676,98 @@ static Optional<Long> parseDataSizeBytes(String value) {
     private void initTable() {
         Preconditions.checkState(params.containsKey("serialized_table"));
         table = PaimonUtils.deserialize(params.get("serialized_table"));
+        table = applyBackendManifestParallelism(table,
+                params.get(PAIMON_OPTION_PREFIX + 
DORIS_MANIFEST_PARALLELISM_CAP),
+                Runtime.getRuntime().availableProcessors());
+        validateSerializedReaderOptions(table);
         paimonAllFieldNames = PaimonUtils.getFieldNames(this.table.rowType());
         if (LOG.isDebugEnabled()) {
             LOG.debug("paimonAllFieldNames:{}", paimonAllFieldNames);
         }
     }
 
+    static Table applyBackendManifestParallelism(
+            Table table, String feParallelismCap, int localCapacity) {
+        int safeParallelism;
+        if (feParallelismCap != null) {
+            safeParallelism = 
parsePositiveManifestParallelism(feParallelismCap);
+            if (safeParallelism <= localCapacity) {
+                return table;
+            }
+            safeParallelism = localCapacity;
+        } else {
+            List<Integer> configuredValues = new ArrayList<>();
+            collectManifestParallelism(table, configuredValues);
+            if (configuredValues.isEmpty()
+                    || configuredValues.stream().noneMatch(value -> value > 
localCapacity)) {
+                return table;
+            }
+            safeParallelism = Math.min(
+                    
configuredValues.stream().mapToInt(Integer::intValue).min().getAsInt(),
+                    localCapacity);
+        }
+        Map<String, String> cap = Collections.singletonMap(
+                CoreOptions.SCAN_MANIFEST_PARALLELISM.key(), 
String.valueOf(safeParallelism));
+        // File-store copies must retain the FE-selected schema while only 
lowering an execution
+        // bound; ordinary copy can re-resolve time travel and undo schema 
pinning.
+        return table instanceof FileStoreTable
+                ? ((FileStoreTable) table).copyWithoutTimeTravel(cap)
+                : table.copy(cap);

Review Comment:
   Valid. Commit bd46e1fa6c1 transports the exact hidden system source and 
rebuilds the same wrapper with copyWithoutTimeTravel on a smaller BE. A real 
PartitionsTable test covers the schema-safe cap.



##########
fe/fe-connector/fe-connector-api/src/main/java/org/apache/doris/connector/api/ConnectorMetadata.java:
##########
@@ -70,6 +71,12 @@ default Optional<ConnectorMvccSnapshot> beginQuerySnapshot(
         return Optional.empty();
     }
 
+    /** Whether relation-level options need the statement's latest snapshot as 
their version fence. */
+    default boolean usesStatementSnapshotForOptions(

Review Comment:
   Valid. Commit bd46e1fa6c1 regenerates the ConnectorMetadata surface 
baseline; ConnectorMetadataSurfaceTest passes.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonReaderOptions.java:
##########
@@ -0,0 +1,331 @@
+// 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.doris.connector.paimon;
+
+import com.google.common.collect.ImmutableSet;
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.options.ConfigOption;
+import org.apache.paimon.options.MemorySize;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.table.DelegatedFileStoreTable;
+import org.apache.paimon.table.FallbackReadFileStoreTable;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.Table;
+import org.apache.paimon.table.system.SystemTableLoader;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.OptionalInt;
+import java.util.Set;
+
+/** Validation shared by catalog-scoped and relation-scoped Paimon reader 
tuning. */
+public final class PaimonReaderOptions {
+    public static final String TABLE_OPTION_PREFIX = "paimon.table-option.";
+    public static final int MIN_READ_BATCH_SIZE = 1;
+    public static final int MAX_READ_BATCH_SIZE = 65536;
+    // Keep catalog replay deterministic while bounding a single option's 
JVM-wide thread impact.
+    public static final int MAX_MANIFEST_PARALLELISM = 256;
+    public static final long MIN_ASYNC_THRESHOLD_BYTES = 1024L * 1024L;
+    public static final long MAX_ASYNC_THRESHOLD_BYTES = 1024L * 1024L * 1024L;
+
+    // Keep this list to batch-read controls consumed by Doris' Paimon scan 
path. Context selectors,
+    // streaming-source settings, storage layout, and write options are unsafe 
after schema binding.
+    private static final Set<String> SUPPORTED_OPTIONS = ImmutableSet.of(
+            CoreOptions.READ_BATCH_SIZE.key(),
+            CoreOptions.FILE_READER_ASYNC_THRESHOLD.key(),
+            CoreOptions.FILE_INDEX_READ_ENABLED.key(),
+            CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key(),
+            CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key(),
+            CoreOptions.SCAN_MANIFEST_PARALLELISM.key(),
+            CoreOptions.SCAN_PLAN_SORT_PARTITION.key());
+
+    // These settings do not alter the selected snapshot or manifest 
projection, so relation-local
+    // copies can reuse the memoized partition projection while still planning 
splits from the copy.
+    private static final Set<String> METADATA_NEUTRAL_OPTIONS = 
ImmutableSet.of(
+            CoreOptions.READ_BATCH_SIZE.key(),
+            CoreOptions.FILE_READER_ASYNC_THRESHOLD.key(),
+            CoreOptions.FILE_INDEX_READ_ENABLED.key(),
+            CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key(),
+            CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key());
+
+    private PaimonReaderOptions() {
+    }
+
+    public static Set<String> supportedOptions() {
+        return SUPPORTED_OPTIONS;
+    }
+
+    public static Set<String> metadataNeutralOptions() {
+        return METADATA_NEUTRAL_OPTIONS;
+    }
+
+    public static void validate(String key, String value) {
+        if (!SUPPORTED_OPTIONS.contains(key)) {
+            throw new IllegalArgumentException("Unsupported Paimon dynamic 
reader option '" + key
+                    + "'. Supported options are " + SUPPORTED_OPTIONS);
+        }
+
+        if (CoreOptions.READ_BATCH_SIZE.key().equals(key)) {
+            int batchSize = parse(key, value, CoreOptions.READ_BATCH_SIZE);
+            // A zero batch can make Paimon's vectorized reader report success 
without
+            // advancing input; the upper bound also prevents one relation 
from over-allocating.
+            requireRange(key, batchSize, MIN_READ_BATCH_SIZE, 
MAX_READ_BATCH_SIZE);
+        } else if (CoreOptions.FILE_READER_ASYNC_THRESHOLD.key().equals(key)) {
+            MemorySize threshold = parse(key, value, 
CoreOptions.FILE_READER_ASYNC_THRESHOLD);
+            // Bound the trigger on both sides so a query cannot fan out tiny 
async reads or
+            // silently disable asynchronous reading with an effectively 
infinite threshold.
+            requireRange(key, threshold.getBytes(),
+                    MIN_ASYNC_THRESHOLD_BYTES, MAX_ASYNC_THRESHOLD_BYTES);
+        } else if (CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key().equals(key)) {
+            MemorySize targetSize = parse(key, value, 
CoreOptions.SOURCE_SPLIT_TARGET_SIZE);
+            // A split-size option represents byte capacity; non-positive 
values silently defeat
+            // Paimon's bin packing and turn every data file into a separate 
Doris scan range.
+            requireRange(key, targetSize.getBytes(), 1, Long.MAX_VALUE);
+        } else if (CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key().equals(key)) {
+            parse(key, value, CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST);
+        } else if (CoreOptions.SCAN_MANIFEST_PARALLELISM.key().equals(key)) {
+            validateManifestParallelism(value);
+        } else if (CoreOptions.FILE_INDEX_READ_ENABLED.key().equals(key)) {
+            parse(key, value, CoreOptions.FILE_INDEX_READ_ENABLED);
+        } else {
+            parse(key, value, CoreOptions.SCAN_PLAN_SORT_PARTITION);
+        }
+    }
+
+    public static void validateCatalogProperties(Map<String, String> 
properties) {
+        properties.forEach((key, value) -> {
+            if (!key.toLowerCase(Locale.ROOT).startsWith(TABLE_OPTION_PREFIX)) 
{
+                return;
+            }
+            String optionKey = key.substring(TABLE_OPTION_PREFIX.length());
+            if (optionKey.isEmpty()) {
+                throw new IllegalArgumentException(
+                        "Paimon table option name must not be empty after 
prefix " + TABLE_OPTION_PREFIX);
+            }
+            validate(optionKey, value);
+        });
+    }
+
+    public static Map<String, String> compatibleCatalogOptions(Map<String, 
String> properties) {
+        Map<String, String> compatibleOptions = new LinkedHashMap<>();
+        properties.forEach((key, value) -> {
+            if (!key.toLowerCase(Locale.ROOT).startsWith(TABLE_OPTION_PREFIX)) 
{
+                return;
+            }
+            String optionKey = key.substring(TABLE_OPTION_PREFIX.length());
+            try {
+                validate(optionKey, value);
+                compatibleOptions.put(optionKey, value);
+            } catch (IllegalArgumentException ignored) {
+                // Images written before the reader-only allowlist may contain 
arbitrary Paimon
+                // options. Keep the catalog loadable, but never apply an 
unsafe legacy option.
+            }
+        });
+        return Collections.unmodifiableMap(compatibleOptions);
+    }
+
+    public static void validateReaderOptions(Map<String, String> options) {
+        SUPPORTED_OPTIONS.stream()
+                .filter(options::containsKey)
+                .forEach(key -> validate(key, options.get(key)));
+    }
+
+    public static void validateEffectiveTableOptions(Map<String, String> 
options) {
+        validateReaderOptions(options);
+        validateIfPresentForRuntime(options, 
CoreOptions.SCAN_MANIFEST_PARALLELISM.key());
+    }
+
+    public static Map<String, String> runtimeSafeCopyOptions(Table table, 
Map<String, String> copyOptions) {
+        Map<String, String> safeOptions = new LinkedHashMap<>(copyOptions);
+        String key = CoreOptions.SCAN_MANIFEST_PARALLELISM.key();
+        if (safeOptions.containsKey(key)) {
+            String configured = safeOptions.get(key);
+            if (configured == null) {
+                return safeOptions;
+            }
+            validateManifestParallelism(configured);
+            int requested = Integer.parseInt(configured);
+            int localCapacity = Runtime.getRuntime().availableProcessors();
+            if (requested > localCapacity) {
+                safeOptions.put(key, String.valueOf(localCapacity));
+            }
+            return safeOptions;
+        }
+
+        OptionalInt safeParallelism = runtimeSafeManifestParallelism(table);
+        if (!safeParallelism.isPresent()) {
+            return safeOptions;
+        }
+        List<Integer> configuredValues = new ArrayList<>();
+        collectManifestParallelism(table, configuredValues);
+        int localCapacity = Runtime.getRuntime().availableProcessors();
+        if (configuredValues.stream().anyMatch(value -> value > 
localCapacity)) {
+            // Keep persisted semantics stable across heterogeneous FEs, but 
cap the execution copy
+            // conservatively across every nested planner hidden by a wrapper.
+            safeOptions.put(key, String.valueOf(safeParallelism.getAsInt()));
+        }
+        return safeOptions;
+    }
+
+    public static OptionalInt runtimeSafeManifestParallelism(Table table) {
+        List<Integer> configuredValues = new ArrayList<>();
+        collectManifestParallelism(table, configuredValues);
+        if (configuredValues.isEmpty()) {
+            return OptionalInt.empty();
+        }
+        int localCapacity = Runtime.getRuntime().availableProcessors();
+        return OptionalInt.of(Math.min(
+                
configuredValues.stream().mapToInt(Integer::intValue).min().getAsInt(),
+                localCapacity));
+    }
+
+    public static Table runtimeSafeTable(Table table) {
+        Map<String, String> runtimeOptions = runtimeSafeCopyOptions(table, 
Collections.emptyMap());
+        // Catalog handles stay hardware-neutral; every local planning 
consumer receives its own
+        // capped copy before it can resize Paimon's JVM-wide manifest 
executor.
+        return runtimeOptions.isEmpty() ? table : table.copy(runtimeOptions);
+    }
+
+    public static Table runtimeSafeSystemTable(
+            String systemTableType, Table systemTable, Table sourceTable, 
Map<String, String> scanOptions) {
+        Table effectiveSource = runtimeSafeSystemSource(sourceTable, 
scanOptions);
+        validateEffectiveTable(effectiveSource);
+        OptionalInt parallelism = 
runtimeSafeManifestParallelism(effectiveSource);
+        if (!parallelism.isPresent()) {
+            return systemTable;
+        }
+        Map<String, String> cap = Collections.singletonMap(
+                CoreOptions.SCAN_MANIFEST_PARALLELISM.key(),
+                String.valueOf(parallelism.getAsInt()));
+        if (effectiveSource instanceof FileStoreTable) {
+            // Copying a system wrapper replays inherited time-travel options 
and can rewind a
+            // schema-only ALTER; cap the source without time travel, then 
rebuild the same wrapper.
+            FileStoreTable cappedSource = ((FileStoreTable) 
effectiveSource).copyWithoutTimeTravel(cap);
+            Table rebuilt = SystemTableLoader.load(systemTableType, 
cappedSource);

Review Comment:
   Valid. Commit bd46e1fa6c1 peels non-fallback decorators before rebuilding 
the capped read-optimized wrapper and covers the immediate fallback pair.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonReaderOptions.java:
##########
@@ -0,0 +1,331 @@
+// 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.doris.connector.paimon;
+
+import com.google.common.collect.ImmutableSet;
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.options.ConfigOption;
+import org.apache.paimon.options.MemorySize;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.table.DelegatedFileStoreTable;
+import org.apache.paimon.table.FallbackReadFileStoreTable;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.Table;
+import org.apache.paimon.table.system.SystemTableLoader;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.OptionalInt;
+import java.util.Set;
+
+/** Validation shared by catalog-scoped and relation-scoped Paimon reader 
tuning. */
+public final class PaimonReaderOptions {
+    public static final String TABLE_OPTION_PREFIX = "paimon.table-option.";
+    public static final int MIN_READ_BATCH_SIZE = 1;
+    public static final int MAX_READ_BATCH_SIZE = 65536;
+    // Keep catalog replay deterministic while bounding a single option's 
JVM-wide thread impact.
+    public static final int MAX_MANIFEST_PARALLELISM = 256;
+    public static final long MIN_ASYNC_THRESHOLD_BYTES = 1024L * 1024L;
+    public static final long MAX_ASYNC_THRESHOLD_BYTES = 1024L * 1024L * 1024L;
+
+    // Keep this list to batch-read controls consumed by Doris' Paimon scan 
path. Context selectors,
+    // streaming-source settings, storage layout, and write options are unsafe 
after schema binding.
+    private static final Set<String> SUPPORTED_OPTIONS = ImmutableSet.of(
+            CoreOptions.READ_BATCH_SIZE.key(),
+            CoreOptions.FILE_READER_ASYNC_THRESHOLD.key(),
+            CoreOptions.FILE_INDEX_READ_ENABLED.key(),
+            CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key(),
+            CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key(),
+            CoreOptions.SCAN_MANIFEST_PARALLELISM.key(),
+            CoreOptions.SCAN_PLAN_SORT_PARTITION.key());
+
+    // These settings do not alter the selected snapshot or manifest 
projection, so relation-local
+    // copies can reuse the memoized partition projection while still planning 
splits from the copy.
+    private static final Set<String> METADATA_NEUTRAL_OPTIONS = 
ImmutableSet.of(
+            CoreOptions.READ_BATCH_SIZE.key(),
+            CoreOptions.FILE_READER_ASYNC_THRESHOLD.key(),
+            CoreOptions.FILE_INDEX_READ_ENABLED.key(),
+            CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key(),
+            CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key());
+
+    private PaimonReaderOptions() {
+    }
+
+    public static Set<String> supportedOptions() {
+        return SUPPORTED_OPTIONS;
+    }
+
+    public static Set<String> metadataNeutralOptions() {
+        return METADATA_NEUTRAL_OPTIONS;
+    }
+
+    public static void validate(String key, String value) {
+        if (!SUPPORTED_OPTIONS.contains(key)) {
+            throw new IllegalArgumentException("Unsupported Paimon dynamic 
reader option '" + key
+                    + "'. Supported options are " + SUPPORTED_OPTIONS);
+        }
+
+        if (CoreOptions.READ_BATCH_SIZE.key().equals(key)) {
+            int batchSize = parse(key, value, CoreOptions.READ_BATCH_SIZE);
+            // A zero batch can make Paimon's vectorized reader report success 
without
+            // advancing input; the upper bound also prevents one relation 
from over-allocating.
+            requireRange(key, batchSize, MIN_READ_BATCH_SIZE, 
MAX_READ_BATCH_SIZE);
+        } else if (CoreOptions.FILE_READER_ASYNC_THRESHOLD.key().equals(key)) {
+            MemorySize threshold = parse(key, value, 
CoreOptions.FILE_READER_ASYNC_THRESHOLD);
+            // Bound the trigger on both sides so a query cannot fan out tiny 
async reads or
+            // silently disable asynchronous reading with an effectively 
infinite threshold.
+            requireRange(key, threshold.getBytes(),
+                    MIN_ASYNC_THRESHOLD_BYTES, MAX_ASYNC_THRESHOLD_BYTES);
+        } else if (CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key().equals(key)) {
+            MemorySize targetSize = parse(key, value, 
CoreOptions.SOURCE_SPLIT_TARGET_SIZE);
+            // A split-size option represents byte capacity; non-positive 
values silently defeat
+            // Paimon's bin packing and turn every data file into a separate 
Doris scan range.
+            requireRange(key, targetSize.getBytes(), 1, Long.MAX_VALUE);
+        } else if (CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key().equals(key)) {
+            parse(key, value, CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST);
+        } else if (CoreOptions.SCAN_MANIFEST_PARALLELISM.key().equals(key)) {
+            validateManifestParallelism(value);
+        } else if (CoreOptions.FILE_INDEX_READ_ENABLED.key().equals(key)) {
+            parse(key, value, CoreOptions.FILE_INDEX_READ_ENABLED);
+        } else {
+            parse(key, value, CoreOptions.SCAN_PLAN_SORT_PARTITION);
+        }
+    }
+
+    public static void validateCatalogProperties(Map<String, String> 
properties) {
+        properties.forEach((key, value) -> {
+            if (!key.toLowerCase(Locale.ROOT).startsWith(TABLE_OPTION_PREFIX)) 
{
+                return;
+            }
+            String optionKey = key.substring(TABLE_OPTION_PREFIX.length());
+            if (optionKey.isEmpty()) {
+                throw new IllegalArgumentException(
+                        "Paimon table option name must not be empty after 
prefix " + TABLE_OPTION_PREFIX);
+            }
+            validate(optionKey, value);
+        });
+    }
+
+    public static Map<String, String> compatibleCatalogOptions(Map<String, 
String> properties) {
+        Map<String, String> compatibleOptions = new LinkedHashMap<>();
+        properties.forEach((key, value) -> {
+            if (!key.toLowerCase(Locale.ROOT).startsWith(TABLE_OPTION_PREFIX)) 
{
+                return;
+            }
+            String optionKey = key.substring(TABLE_OPTION_PREFIX.length());
+            try {
+                validate(optionKey, value);
+                compatibleOptions.put(optionKey, value);
+            } catch (IllegalArgumentException ignored) {
+                // Images written before the reader-only allowlist may contain 
arbitrary Paimon
+                // options. Keep the catalog loadable, but never apply an 
unsafe legacy option.
+            }
+        });
+        return Collections.unmodifiableMap(compatibleOptions);
+    }
+
+    public static void validateReaderOptions(Map<String, String> options) {
+        SUPPORTED_OPTIONS.stream()
+                .filter(options::containsKey)
+                .forEach(key -> validate(key, options.get(key)));
+    }
+
+    public static void validateEffectiveTableOptions(Map<String, String> 
options) {
+        validateReaderOptions(options);
+        validateIfPresentForRuntime(options, 
CoreOptions.SCAN_MANIFEST_PARALLELISM.key());
+    }
+
+    public static Map<String, String> runtimeSafeCopyOptions(Table table, 
Map<String, String> copyOptions) {
+        Map<String, String> safeOptions = new LinkedHashMap<>(copyOptions);
+        String key = CoreOptions.SCAN_MANIFEST_PARALLELISM.key();
+        if (safeOptions.containsKey(key)) {
+            String configured = safeOptions.get(key);
+            if (configured == null) {
+                return safeOptions;
+            }
+            validateManifestParallelism(configured);
+            int requested = Integer.parseInt(configured);
+            int localCapacity = Runtime.getRuntime().availableProcessors();
+            if (requested > localCapacity) {
+                safeOptions.put(key, String.valueOf(localCapacity));
+            }
+            return safeOptions;
+        }
+
+        OptionalInt safeParallelism = runtimeSafeManifestParallelism(table);
+        if (!safeParallelism.isPresent()) {
+            return safeOptions;
+        }
+        List<Integer> configuredValues = new ArrayList<>();
+        collectManifestParallelism(table, configuredValues);
+        int localCapacity = Runtime.getRuntime().availableProcessors();
+        if (configuredValues.stream().anyMatch(value -> value > 
localCapacity)) {
+            // Keep persisted semantics stable across heterogeneous FEs, but 
cap the execution copy
+            // conservatively across every nested planner hidden by a wrapper.
+            safeOptions.put(key, String.valueOf(safeParallelism.getAsInt()));
+        }
+        return safeOptions;
+    }
+
+    public static OptionalInt runtimeSafeManifestParallelism(Table table) {
+        List<Integer> configuredValues = new ArrayList<>();
+        collectManifestParallelism(table, configuredValues);
+        if (configuredValues.isEmpty()) {
+            return OptionalInt.empty();
+        }
+        int localCapacity = Runtime.getRuntime().availableProcessors();
+        return OptionalInt.of(Math.min(
+                
configuredValues.stream().mapToInt(Integer::intValue).min().getAsInt(),
+                localCapacity));
+    }
+
+    public static Table runtimeSafeTable(Table table) {
+        Map<String, String> runtimeOptions = runtimeSafeCopyOptions(table, 
Collections.emptyMap());
+        // Catalog handles stay hardware-neutral; every local planning 
consumer receives its own
+        // capped copy before it can resize Paimon's JVM-wide manifest 
executor.
+        return runtimeOptions.isEmpty() ? table : table.copy(runtimeOptions);
+    }
+
+    public static Table runtimeSafeSystemTable(
+            String systemTableType, Table systemTable, Table sourceTable, 
Map<String, String> scanOptions) {
+        Table effectiveSource = runtimeSafeSystemSource(sourceTable, 
scanOptions);
+        validateEffectiveTable(effectiveSource);
+        OptionalInt parallelism = 
runtimeSafeManifestParallelism(effectiveSource);
+        if (!parallelism.isPresent()) {
+            return systemTable;
+        }
+        Map<String, String> cap = Collections.singletonMap(
+                CoreOptions.SCAN_MANIFEST_PARALLELISM.key(),
+                String.valueOf(parallelism.getAsInt()));
+        if (effectiveSource instanceof FileStoreTable) {
+            // Copying a system wrapper replays inherited time-travel options 
and can rewind a
+            // schema-only ALTER; cap the source without time travel, then 
rebuild the same wrapper.
+            FileStoreTable cappedSource = ((FileStoreTable) 
effectiveSource).copyWithoutTimeTravel(cap);
+            Table rebuilt = SystemTableLoader.load(systemTableType, 
cappedSource);
+            if (rebuilt == null) {
+                throw new IllegalArgumentException("Unsupported Paimon system 
table '"
+                        + systemTableType + "'");
+            }
+            return rebuilt;
+        }
+        return systemTable.copy(cap);
+    }
+
+    public static Table runtimeSafeSystemSource(Table sourceTable, Map<String, 
String> scanOptions) {
+        return PaimonScanParams.isOptionsPin(scanOptions)
+                ? PaimonScanParams.applyOptions(sourceTable, scanOptions)
+                : runtimeSafeTable(sourceTable);

Review Comment:
   Valid. Commit bd46e1fa6c1 reapplies reset-aware incremental options to the 
exact source before capping and rebuilding the system wrapper.



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


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to