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]