github-actions[bot] commented on code in PR #66247:
URL: https://github.com/apache/doris/pull/66247#discussion_r3698709652


##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonConnectorMetadata.java:
##########
@@ -749,8 +768,12 @@ public Optional<ConnectorMvccSnapshot> resolveTimeTravel(
                             
.properties(PaimonScanParams.markAsOptions(resolved))
                             .build());
                 }
-                long pinnedId = pinnedSnapshotId(table, resolved);
-                long schemaId = pinnedId < 0
+                long pinnedId = usesStatementFence
+                        ? spec.getLatestSnapshotFence().getAsLong()
+                        : pinnedSnapshotId(table, resolved);
+                // The statement fence pins data visibility, not schema time 
travel. Planning-only
+                // aliases must retain the latest-schema projection used by 
the plain relation.
+                long schemaId = usesStatementFence || pinnedId < 0

Review Comment:
   [P1] Propagate this latest-schema sentinel to system-table binding. The 
normal `getTableSchema` arm honors `schemaId=-1`, but its system-table arm 
bypasses that check and `resolveSystemTableAt` always applies the marked 
OPTIONS map; the synthetic `scan.snapshot-id=S` therefore still rewinds the 
view's row type. After a schema-only add/rename following S, 
`$ro@options('scan.manifest.parallelism'='1')` binds S-era fields while the 
runtime-safe wrapper keeps the current schema. Please retain/honor the 
statement-fence provenance in the system path and cover both planning-only and 
explicit selectors.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonReaderOptions.java:
##########
@@ -0,0 +1,347 @@
+// 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.privilege.PrivilegedFileStoreTable;
+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.Collections;
+import java.util.LinkedHashMap;
+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 OptionalInt backendManifestParallelismCap(Table table) {
+        // The transport value is the execution ceiling, not the smallest 
branch preference;
+        // the BE preserves each branch's lower value while independently 
capping larger siblings.
+        return OptionalInt.of(Math.min(
+                Runtime.getRuntime().availableProcessors(), 
MAX_MANIFEST_PARALLELISM));
+    }
+
+    public static Table runtimeSafeTable(Table table) {
+        return runtimeSafeTable(table, 
Runtime.getRuntime().availableProcessors());
+    }
+
+    static Table runtimeSafeTable(Table table, int localCapacity) {
+        if (localCapacity < 1) {
+            throw new IllegalArgumentException("Paimon planning capacity must 
be positive.");
+        }
+        int safeBound = Math.min(localCapacity, MAX_MANIFEST_PARALLELISM);
+        return normalizeManifestParallelism(table, safeBound, localCapacity > 
safeBound);
+    }
+
+    public static Table runtimeSafeSystemTable(
+            String systemTableType, Table systemTable, Table sourceTable, 
Map<String, String> scanOptions) {
+        Table effectiveSource = runtimeSafeSystemSource(sourceTable, 
scanOptions);
+        validateEffectiveTable(effectiveSource);
+        if (effectiveSource instanceof FileStoreTable) {
+            // Paimon dispatches fallback reads only when the fallback pair is 
the system wrapper's
+            // immediate child; a privilege decorator must not hide that pair 
during this rebuild.
+            FileStoreTable systemSource = 
PaimonTableDecorators.unwrapToFallbackOrBase(
+                    (FileStoreTable) effectiveSource);
+            Table rebuilt = SystemTableLoader.load(systemTableType, 
systemSource);
+            if (rebuilt == null) {
+                throw new IllegalArgumentException("Unsupported Paimon system 
table '"
+                        + systemTableType + "'");
+            }
+            return rebuilt;
+        }
+        return runtimeSafeTable(systemTable);
+    }
+
+    public static Table runtimeSafeSystemSource(Table sourceTable, Map<String, 
String> scanOptions) {
+        if (PaimonScanParams.isOptionsPin(scanOptions)) {
+            if (sourceTable instanceof FileStoreTable) {
+                return PaimonScanParams.applyOptionsWithoutTimeTravel(

Review Comment:
   [P1] Preserve historical schemas for explicit system-table OPTIONS. This 
branch uses `copyWithoutTimeTravel` for every marked OPTIONS pin, but a scan 
handle's `systemTableSource` is freshly resolved at the latest generation. 
After `old` is renamed to `new`, 
`$ro`/`$audit_log@options('scan.snapshot-id'='S')` binds `old` from S, then 
this rebuild exposes the latest row type to `planScanInternal`; its OPTIONS 
projection maps `old` to -1 and aborts before the backend rebuild can restore 
S. Planning-only OPTIONS do need latest-schema semantics, but explicit 
selectors need the selected schema. Please retain that provenance/source 
generation and cover split planning for a historical system view across a 
rename.



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